"""Example Kafka Producer in Flask Application.""" from contextlib import contextmanager import json import logging from random import randint from flask import Flask, request from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer # config TOPIC = 'event.examplePythonProducer' BOOTSTRAP_SERVERS = ( 'b-2.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094,' 'b-3.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094,' 'b-1.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094' ) PROTOCOL = 'SSL' # or PLAINTEXT if local non-ssl # connectors logger = logging.getLogger('kafka_producer_logger') # Main Application app = Flask(__name__) @app.route("/") def hello_world(): return """
If you click the "Submit" button, the form-data will be sent to a handler called "/produce".
""" @app.route("/produce", methods=['POST']) def new_record(): def _on_delivery(err, msg): """Handle produce result.""" if err is not None: raise Exception(err) logger.info('SUCCESS!') logger.info('KEY: ', msg.key()) logger.info('TOPIC: ', msg.topic()) logger.info('PARTITION:', msg.partition()) logger.info('OFFSET: ', msg.offset()) producer = EventProducer( BOOTSTRAP_SERVERS, StringSerializer(), StringSerializer(), PROTOCOL) event = { 'person_id': randint(1, 1000), 'first_name': request.form['fname'], 'last_name': request.form['lname']} producer.produce( topic=TOPIC, event_key=str(randint(0, 1000)), event_value=json.dumps(event), callback=_on_delivery) return {'emitted_record': event} if __name__ == '__main__': app.run( host='localhost', port=8000, debug=True)