【问题标题】:I have an issue posting JSON message on Kafka using PyKafka我在使用 PyKafka 在 Kafka 上发布 JSON 消息时遇到问题
【发布时间】:2021-08-10 23:45:33
【问题描述】:

我正在使用 pykafka 库在 Kafka 上发布消息。我的数据集是 JSON

{"user": "jpoole", "created_at_unixtime": 1440407147.033846, "id": 3600730356622213650, "text": "Techical support for my new computer as A+, thank you @fudgemart", "created_at": "Mon Aug 24 05:05:47 +0000 2015"} 
]

我的要求是生成 2 条 kafka 消息,使用 PyKafka 为上面的每个 JSON 字符串生成 1 条。到目前为止,我已经尝试了以下方法。

from pykafka import KafkaClient

client = KafkaClient(hosts="127.0.0.1:9092")
topic = client.topics['test']

with open('./tweets.json') as f:
    dataItems =json.load(f)

s=json.dumps(dataItems).encode('utf-8')


with topic.get_sync_producer() as producer:
    for data in s:
        producer.produce(data)

我已将 JSON 加载到文件中(我最初的要求)。上面的代码有效,但它没有将第一个 JSON 字符串作为一个整体,而是将字符串中的每个字符都作为一条消息。

我的要求是将每个 JSON 字符串作为单独的 Kafka 消息发布。

Message 1
{"user": "jpoole", "created_at_unixtime": 1448221456.6646008, "id": 3731785240073317438, "text": "Glad I bought my electronics from @fudgemart", "created_at": "Sun Nov 22 14:44:16 +0000 2015"}

Message 2
{"user": "jpoole", "created_at_unixtime": 1440407147.033846, "id": 3600730356622213650, "text": "Techical support for my new computer as A+, thank you @fudgemart", "created_at": "Mon Aug 24 05:05:47 +0000 2015"}

谢谢

【问题讨论】:

    标签: json pykafka


    【解决方案1】:

    没有任何日志很难说,但我认为问题可能出在for data in s: 这一行。

    JSON.dumps() 产生一个字符串,所以s 是一个字符串,data 是一个字符。 所以问题是你要分别生成每个字符。

    【讨论】:

    • 添加了 json.dumps() 以将字典转换为字节。 producer.produce() 方法需要一个字节。如果我不转换这是错误,我得到“TypeError: ("Producer.produce 接受一个字节对象作为消息,但它得到了 '%s'", )"
    • 据我了解,JSON.dumps 返回一个字符串:docs.python.org/3/library/json.html#json.dump。因此,在这里,您尝试直接发送对象(这是行不通的,如错误消息所示)。而且您还尝试发送单个字符串。但是据我了解,您还没有尝试发送实际的字符串:即使用JSON.dumps,然后直接返回(不逐个字符)。如果我的理解有误,请纠正我。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-17
    • 1970-01-01
    • 1970-01-01
    • 2021-10-29
    • 1970-01-01
    相关资源
    最近更新 更多