【发布时间】: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