#!/usr/bin/env python3 from kafka import KafkaConsumer import boto3 import json import os import sys from datetime import datetime # Configuratie 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("=" * 60) print(f"Kafka Consumer starting...") print(f"Kafka broker: {KAFKA_BROKER}") print(f"Kafka topic: {KAFKA_TOPIC}") print(f"S3 endpoint: {S3_ENDPOINT}") print(f"S3 bucket: {S3_BUCKET}") print("=" * 60) # S3 client 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 ) # Kafka consumer try: consumer = KafkaConsumer( KAFKA_TOPIC, bootstrap_servers=[KAFKA_BROKER], auto_offset_reset='earliest', enable_auto_commit=True, value_deserializer=lambda x: json.loads(x.decode('utf-8')) ) print("Connected to Kafka successfully!") except Exception as e: print(f"Failed to connect to Kafka: {e}") sys.exit(1) print(f"Listening for messages on topic '{KAFKA_TOPIC}'...") print("-" * 60) message_count = 0 for message in consumer: message_count += 1 timestamp = datetime.now().strftime('%Y%m%d_%H%M%S_%f') key = f"kafka/{KAFKA_TOPIC}/message_{timestamp}.json" try: s3.put_object( Bucket=S3_BUCKET, Key=key, Body=json.dumps(message.value, indent=2) ) print(f"[{message_count}] Written to s3://{S3_BUCKET}/{key}") print(f" Data: {json.dumps(message.value)[:100]}...") except Exception as e: print(f"Error writing to S3: {e}")