Files

49 lines
1.2 KiB
Python
Raw Permalink Normal View History

from kafka import KafkaConsumer
import boto3
import json
from datetime import datetime
KAFKA_BROKER = "10.0.21.36:9092"
KAFKA_TOPIC = "test-lakehouse"
S3_ENDPOINT = "http://10.0.20.111:9020"
S3_BUCKET = "data"
S3_ACCESS_KEY = "REDACTED"
S3_SECRET_KEY = "REDACTED"
print("Starting consumer...")
print(f"Kafka: {KAFKA_BROKER}, Topic: {KAFKA_TOPIC}")
s3 = boto3.client(
's3',
endpoint_url=S3_ENDPOINT,
aws_access_key_id=S3_ACCESS_KEY,
aws_secret_access_key=S3_SECRET_KEY,
use_ssl=False,
verify=False
)
def deserialize(value):
if value is None or len(value) == 0:
return None
try:
return json.loads(value.decode('utf-8'))
except:
print(f"Skipping non-JSON")
return None
consumer = KafkaConsumer(
KAFKA_TOPIC,
bootstrap_servers=[KAFKA_BROKER],
auto_offset_reset='latest',
value_deserializer=deserialize
)
print("Waiting for messages...")
for msg in consumer:
if msg.value is None:
continue
ts = datetime.now().strftime('%Y%m%d_%H%M%S_%f')
key = f"kafka/{KAFKA_TOPIC}/msg_{ts}.json"
s3.put_object(Bucket=S3_BUCKET, Key=key, Body=json.dumps(msg.value))
print(f"Written: {key}")