Files
Lakehouse/config/spark-jobs/s3_consumer_realtime.py
T

46 lines
1.1 KiB
Python
Raw Normal View History

from kafka import KafkaConsumer
import boto3
import json
from datetime import datetime
KAFKA_BROKER = "10.0.21.36:9092"
S3_ENDPOINT = "http://10.0.20.111:9020"
S3_BUCKET = "data"
S3_ACCESS_KEY = "REDACTED"
S3_SECRET_KEY = "REDACTED"
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
)
# Consumer die naar ALLE topics luistert (alleen nieuwe berichten)
consumer = KafkaConsumer(
bootstrap_servers=[KAFKA_BROKER],
auto_offset_reset='latest'
)
# Subscribe op alle topics
consumer.subscribe(pattern='.*')
print("Listening to ALL topics...")
print("Waiting for new messages...")
msg_count = 0
for msg in consumer:
msg_count += 1
timestamp = datetime.now().strftime('%Y%m%d_%H%M%S_%f')
key = f"kafka/{msg.topic}/message_{timestamp}.json"
try:
data = json.loads(msg.value.decode('utf-8'))
except:
data = {"raw": msg.value.decode('utf-8')}
s3.put_object(Bucket=S3_BUCKET, Key=key, Body=json.dumps(data, indent=2))
print(f"[{msg_count}] {msg.topic} -> s3://{S3_BUCKET}/{key}")