基于Twisted框架的Python应用对接Kafka服务器的问题求助
基于Twisted框架的Python应用对接Kafka服务器的问题求助
我有一个用Twisted框架开发的应用,它通过输出插件将各类数据(主要是JSON格式的日志条目)发送给不同的日志服务器。现在我想扩展它的功能,让它也能把数据发送到Kafka服务器,但遇到了一个棘手的问题,不知道该怎么解决。
如果我用Python的kafka-python模块以常规方式直接向Kafka服务器发送数据,一切都能正常工作,代码示例如下:
from json import dumps from kafka import KafkaProducer site = '<KAFKA SERVER>' port = 9092 username = '<USERNAME>' password = '<PASSWORD>' topic = 'test' producer = KafkaProducer( bootstrap_servers='{}:{}'.format(site, port), value_serializer=lambda v: bytes(str(dumps(v)).encode('utf-8')), sasl_mechanism='SCRAM-SHA-256', security_protocol='SASL_SSL', sasl_plain_username=username, sasl_plain_password=password ) event = { 'message': 'Test message' } try: producer.send(topic, value=event) producer.flush() print("消息发送成功") except Exception as e: print(f"发送失败: {str(e)}")
但一旦尝试把这套发送逻辑整合到Twisted应用中就出现了问题,我暂时没找到问题根源,希望能得到大家的帮助。
备注:内容来源于stack exchange,提问作者bontchev
相关产品推荐
相关产品推荐

