/
LarsFillmore
/
python_basic-flask-kafka
Обзор
Документация
Войти
/
LarsFillmore
/
python_basic-flask-kafka
Код
Запросы
0
Задачи
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
flask-kafka/basic_producer.py
95 строк
3 KB
SoullessSoldier
worked version 0.0.1
05 сен 2024, 20:44
05 сен 2024, 20:44
7505cd4
Код
Авторство
О чём код?
import json import threading from flask import Flask, request, jsonify from confluent_kafka import Producer, KafkaException import time import logging from apscheduler.schedulers.background import BackgroundScheduler app = Flask(__name__) KAFKA_CONFIG = { 'bootstrap.servers': 'kafka-0:9094', 'client.id': 'python-producer', 'acks': 'all', # Request acknowledgments to ensure messages are properly written 'retries': 3, # Number of retries before giving up 'batch.size': 16384, # Batch size for sending messages 'linger.ms': 100, # Adding delay to aggregate messages into larger batches 'compression.type': 'gzip' # Enable compression to save bandwidth } KAFKA_TOPIC = 'test1' producer = Producer(KAFKA_CONFIG) buffer_lock = threading.Lock() flush_interval = 10 buffer = [] chunk_size = 100 # Number of messages to send at a time logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') logger = logging.getLogger(__name__) def delivery_report(err, msg): if err is not None: logger.error(f"Message delivery failed: {err}") else: logger.info(f"Message delivered to {msg.topic()} [{msg.partition()}]") def wait_for_connection(producer, timeout=60): start_time = time.time() while True: try: producer.poll(1) if producer.list_topics(timeout=5): logger.info("Successfully connected to Kafka") return True except KafkaException as e: logger.error(f"Error during connection check: {e}") if time.time() - start_time > timeout: logger.error("Timeout reached, could not connect to Kafka") return False time.sleep(1) def flush_buffer(): logger.info('flush_buffer tick') with buffer_lock: logger.info('flush_buffer tock, len of buffer %d', len(buffer)) while buffer: chunk = buffer[:chunk_size] buffer[:] = buffer[chunk_size:] logger.info('Processing chunk of size %d', len(chunk)) if wait_for_connection(producer): try: for message in chunk: logger.info(message.get('user_name', 'Unknown user')) key = message.get('user_name', 'Unknown')[0] producer.produce(KAFKA_TOPIC, value=json.dumps(message).encode(), key=key, callback=delivery_report) producer.flush() except KafkaException as e: logger.error(f"Error sending messages in chunk: {e}") else: logger.error("Could not connect to Kafka, exiting.") break @app.route('/post_data', methods=['POST']) def post_data(): data = request.json logger.info('Received data: %s', data) with buffer_lock: buffer.append(data) return jsonify({"status": "received"}) if __name__ == '__main__': scheduler = BackgroundScheduler() scheduler.add_job(func=flush_buffer, trigger="interval", seconds=flush_interval) scheduler.start() app.run(debug=False, host='0.0.0.0', port=5000)