ahj*_*ahj 1 python producer message-bus apache-kafka kafka-python
如何在pykafka主题的特定分区上发布消息.在下面的代码片段中,测试主题有四个分区,我打算在其中一个分区中编写每个消息,但显然它不是那样工作的.
from pykafka import KafkaClient
import logging
logging.basicConfig()
client = KafkaClient(hosts='localhost:9092')
print client.topics
topic = client.topics['test']
with topic.get_producer() as producer:
for i in range(4):
producer.produce('another test message ' + str(i ** 2), partition_key='{}'.format(0))
Run Code Online (Sandbox Code Playgroud)
关键是什么决定"哪个分区"的消息将会在结束了.
如果你不提供钥匙,然后把卡夫卡中的消息循环方式,其中每个分区大致得到的消息是相同的.
如果提供密钥,则Kafka会计算哈希并将消息放入生成的分区中.您无法完全控制将使用哪个特定分区,只是相同的密钥将始终在同一分区中.
添加消息密钥通常用于保证某些消息子集的排序.例如,假设您拥有user和transaction实体,并且您希望按顺序处理与同一用户相关的所有交易.您可以通过使用userId消息键来实现这一点.
分区之间没有协调(太慢),因此在使用多个分区时没有总排序.只要将消息全部放在同一个分区中,就可以保证消息的生成顺序与消息生成顺序相同.
也许我应该首先问你的用例,然后再写这一切:)
| 归档时间: |
|
| 查看次数: |
851 次 |
| 最近记录: |