【问题标题】:how i can use multiple consumer in python-kafka?我如何在 python-kafka 中使用多个消费者?
【发布时间】:2019-08-17 19:55:08
【问题描述】:

我看到了文档。 https://kafka-python.readthedocs.io/en/master/usage.html

我尝试了多个消费者。 但它不工作? 以及它是如何工作的?

怎么了?

consumer1 = KafkaConsumer(
     bootstrap_servers=['localhost:9092'],
     #auto_offset_reset='earliest' , # 'earliest',
     #enable_auto_commit= False ,
     group_id='sr',
     value_deserializer=lambda x: loads(x.decode('utf-8')))
consumer1.subscribe("numtest")
consumer2 = KafkaConsumer(
     bootstrap_servers=['localhost:9092'],
     #auto_offset_reset='earliest' , # 'earliest',
     #enable_auto_commit= False ,
     group_id='sr',
     value_deserializer=lambda x: loads(x.decode('utf-8')))
consumer2.subscribe("numtest")
producer = KafkaProducer(bootstrap_servers=['localhost:9092'],
                         key_serializer = None ,                         
                         value_serializer=lambda x: 
                         dumps(x).encode('utf-8'))
def on_send_success(record_metadata):
    print("topic : {} , partition : {} , offset : {}".\
          format( record_metadata.topic , record_metadata.partition , record_metadata.offset))

for msg , message in zip(consumer1 , consumer2) :
    print("="*50)
    print ("topic=%s partition=%d offset=%d: key=%s value=%s" % (message.topic, message.partition,
                                          message.offset, message.key,
                                          message.value))
    print ("topic=%s partition=%d offset=%d: key=%s value=%s" % (msg.topic, msg.partition,
                                          msg.offset, msg.key,
                                          msg.value))
    key = str(message.offset) + " " + str(msg.offset)
    producer.send('output', value= {  key : key }  ).add_callback(on_send_success)
    print("="*50)

其实我想做的是我想让两台电脑为模型的展示做一个ex输入并合并两个结果。

相反,我在合并时必须保持相同的偏移量。

之后,我想发送结果输出主题。

我在outputtopic中的预期结果是这样的:

{ offset = 1 : offest = 1} , { offset = 2 : offest = 2 } , ....

请帮助我!我解决不了

【问题讨论】:

    标签: python apache-kafka


    【解决方案1】:

    您能否详细说明一下您的预期结果是什么?

    如果你想这样做,

    1. 创建两个不同的消费者,其中消费者 1 和消费者 2 获得相同的消息
      • 为此,两个消费者的组 id 应该不同,尝试使用 'sr1' 和 'sr2' 作为组 id。试试下面的代码
    consumer1 = KafkaConsumer(
         bootstrap_servers=['localhost:9092'],
         #auto_offset_reset='earliest' , # 'earliest',
         #enable_auto_commit= False ,
         group_id='sr',
         value_deserializer=lambda x: loads(x.decode('utf-8')))
    consumer1.subscribe("numtest")
    consumer2 = KafkaConsumer(
         bootstrap_servers=['localhost:9092'],
         #auto_offset_reset='earliest' , # 'earliest',
         #enable_auto_commit= False ,
         group_id='sr',
         value_deserializer=lambda x: loads(x.decode('utf-8')))
    consumer2.subscribe("numtest")
    
    1. 创建一个消费者组,其中消费者 1 获取部分消息,消费者 2 获取其余消息
      • 如果您尝试实现这一点,我认为您当前的代码也可以正常工作,但请记住,您可能不想这样做,因为不同的消费者通常分布在不同的进程中。

    【讨论】:

    • 非常感谢您的快速回复!事实上,我想要做的是我想要两台计算机为模型的演示做一个 ex 输入并合并两个结果。相反,当我合并时,我必须保持相同的偏移量。之后,我想发送结果输出主题。我的预期结果在输出主题中是这样的: { offset = 1 : offest = 1} , { offset = 2 : offest = 2 } , .... 请帮帮我!我解决不了
    • 他们来自两个不同的主题吗?在这种情况下,您是否尝试执行类似 Join (dzone.com/articles/join-semantics-in-kafka-streams) 的操作?
    • 我想在使用相同主题“numest”的同时将主题和每个消费者之间的处理值组合在一起,同时使用相同的偏移量。此外,我想对每个消费者使用多处理。例如 {"same offset" , "consumer1_reuslt" + "consumer_result2"} {"0" , "C1 , 0.78 , C2 , 0.1"} {"1" , "C1 , 0.85 , C2 , 0.57"}
    猜你喜欢
    • 2015-09-03
    • 2019-07-22
    • 2019-04-07
    • 1970-01-01
    • 1970-01-01
    • 2020-05-26
    • 2019-11-18
    • 2016-07-16
    相关资源
    最近更新 更多