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}")