【问题标题】:Poll several messages from Kafka轮询来自 Kafka 的几条消息
【发布时间】:2019-09-01 06:22:47
【问题描述】:

我正在使用 confluent_kafka 包来处理 Kafka。 我以这种方式创建主题:

from confluent_kafka import avro
from confluent_kafka.avro import AvroProducer

def my_producer():
    bootstrap_servers=['my_adress.com:9092',
                    'my_adress.com:9092']

    value_schema = avro.load('/home/ValueSchema.avsc')

    avroProducer = AvroProducer({
        'bootstrap.servers': bootstrap_servers[0]+','+bootstrap_servers[1],
        'schema.registry.url':'http://my_adress.com:8081',
        },
        default_value_schema=value_schema
        )

    for i in range(0, 25000):
        value = {"name":"Yuva","favorite_number":10,"favorite_color":"green","age":i*2}
        avroProducer.produce(topic='my_topik14', value=value)
        avroProducer.flush(0)
    print('Finished!')


if __name__ == '__main__':
    my_producer()

它有效。 (这会得到 24820 条消息,而不是 25000 条......) 我们可以检查一下:

kafka-run-class kafka.tools.GetOffsetShell --broker-list my_adress.com:9092 --topic my_topik14
my_topik14:0:24819

现在我想消费:

from confluent_kafka import KafkaError
from confluent_kafka.avro import AvroConsumer
from confluent_kafka.avro.serializer import SerializerError

bootstrap_servers=['my_adress.com:9092',
                   'my_adress.com:9092']
c = AvroConsumer(
    {'bootstrap.servers': bootstrap_servers[0]+','+bootstrap_servers[1],
     'group.id': 'avroneversleeps',
     'schema.registry.url': 'http://my_adress.com:8081',
     'api.version.request': True,
     'fetch.min.bytes': 100000,
     'consume.callback.max.messages':1000,
     'batch.num.messages':2
     })
c.subscribe(['my_topik14'])
running = True

while running:
    msg = None
    try:
        msg = c.poll(0.1)
        if msg:
            if not msg.error():
                print(msg.value())
                c.commit(msg)
            elif msg.error().code() != KafkaError._PARTITION_EOF:
                print(msg.error())
                running = False
        else:
            print("No Message!! Happily trying again!!")
    except SerializerError as e:
        print("Message deserialization failed for %s: %s" % (msg, e))
        running = False
c.commit()
c.close()

但是有一个问题: 我只是一条一条地阅读消息。 我的问题是如何阅读批量消息? 我在 Consumer config 中尝试了不同的参数,但它们没有改变任何东西!


我还在 SO 上找到 this question 并尝试了相同的参数 - 它仍然不起作用。

也读这个。但这与上一个链接相反...

【问题讨论】:

  • 你可以试试这个来计算从bin/kafka-console-consumer.sh --new-consumer --bootstrap-server my_adress.com:9092 --topic my_topik14 --from-beginning | wc 开始的消息数量,看看是否返回25K消息?
  • @St1id3r 感谢您的回复。我不知道如何修改您的命令以使其工作!如果我只是复制它,我会得到bin/kafka-console-consumer.sh: No such file or directory。如果我在命令中添加参数--from-beginning,我会得到`from-beginning is not a known option`
  • 这很奇怪,您应该在 bin 中拥有该实用程序。文档链接:kafka.apache.org/quickstart#quickstart_consume
  • 你可以试试kafka-console-consumer.sh --new-consumer --bootstrap-server my_adress.com:9092 --topic my_topik14 --from-beginning | wc。从上面的命令kafka-run-class kafka.tools.GetOffsetShell --broker-list my_adress.com:9092 --topic my_topik14 my_topik14:0:24819 看来,您的路径上有这些实用程序
  • 我的错,在上面的命令中也省略了 .sh。

标签: python apache-kafka kafka-consumer-api avro confluent-platform


【解决方案1】:

【讨论】:

  • 感谢您的回复。您的问题有一个小错误:AvroConsumer 没有 consume 方法 - 第二个链接地址与第一个链接地址相同。
  • 但是 github issue 中的代码让我想到了 AvroConsumer 类的消费实现。我会在这里展示。
  • 请看我的回答!
【解决方案2】:

AvroConsumer 没有 consume 方法。但是很容易让我自己实现这个方法,就像在 Consume 类(AvroConsumer 的父级)中一样。 代码如下:

def consume_batch(self, num_messages=1, timeout=None):
    """
    This is an overriden method from confluent_kafka.Consumer class. This handles batch of message
    deserialization using avro schema

    :param int num_messages: number of messages to read in one batch (default=1)
    :param float timeout: Poll timeout in seconds (default: indefinite)
    :returns: list of messages objects with deserialized key and value as dict objects
    :rtype: Message
    """
    messages_out = []
    if timeout is None:
        timeout = -1
    messages = super(AvroConsumer, self).consume(num_messages=num_messages, timeout=timeout)
    if messages is None:
        return None
    else:
        for m in messages:
            if not m.value() and not m.key():
                return messages
            if not m.error():
                if m.value() is not None:
                    decoded_value = self._serializer.decode_message(m.value())
                    m.set_value(decoded_value)
                if m.key() is not None:
                    decoded_key = self._serializer.decode_message(m.key())
                    m.set_key(decoded_key)
                messages_out.append(m)
    #print(len(message))
    return messages_out

但是在那之后我们运行测试并且这个方法没有任何性能提升。所以看起来只是为了更好的可用性。或者我需要做一些额外的工作来序列化不是单个消息,而是整个批次。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-03-19
    • 2018-06-04
    • 1970-01-01
    • 2021-09-09
    • 2018-09-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多