【问题标题】:Python produce to different Kafka partitionPython产生到不同的Kafka分区
【发布时间】:2020-09-10 17:09:25
【问题描述】:

我正在尝试通过经典的 Twitter 流示例来学习 Kafka。我正在尝试使用我的生产者将基于 2 个过滤器的 twitter 数据流式传输到同一主题的不同分区。例如,track='Google' 到一个分区和 track='Apple' 到另一个分区的 twitter 数据。

class Producer(StreamListener):
    def __init__(self, producer):
        self.producer = producer

    def on_data(self, data):
        self.producer.send(topic_name, value=data)
        return True

    def on_error(self, error):
        print(error)


twitter_stream = Stream(auth, Producer(producer))
twitter_stream.filter(track=["Google"])

如何添加另一个轨道并将该数据流式传输到另一个分区。

同样,我如何让我的消费者从特定分区消费。

consumer = KafkaConsumer(
    topic_name,
     bootstrap_servers=['localhost:9092'],
     auto_offset_reset='latest',
     enable_auto_commit=True,
     auto_commit_interval_ms =  5000,
     max_poll_records = 100,
     value_deserializer=lambda x: json.loads(x.decode('utf-8')))

【问题讨论】:

    标签: python apache-kafka twitter-streaming-api


    【解决方案1】:

    经过一番研究,我能够解决这个问题:

    在生产者端,指定分区:

    self.producer.send(topic_name, value=data,partition=0)
    

    在消费者方面,

    consumer = KafkaConsumer(
           bootstrap_servers=['localhost:9092'],
         auto_offset_reset='latest',
         enable_auto_commit=True,
         auto_commit_interval_ms =  5000,
         max_poll_records = 100,
         value_deserializer=lambda x: json.loads(x.decode('utf-8')))
    consumer.assign([TopicPartition('trial', 0)])
    

    【讨论】:

      【解决方案2】:

      Kafka 根据消息的键对数据进行分区。在您给定的代码中,您仅将value 传递给生产者消息,因此密钥将为空,因此将在所有分区之间循环。

      请参阅 Kafka 库的文档,了解如何为每条消息指定密钥

      【讨论】:

        猜你喜欢
        • 2019-10-05
        • 2018-12-15
        • 1970-01-01
        • 2023-03-26
        • 2017-04-28
        • 2020-12-29
        • 2019-02-03
        • 2020-01-12
        • 2017-10-18
        相关资源
        最近更新 更多