You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka消费者反序列化失败:未知协议ID问题求助

问题原因

Quix Streams默认使用自定义二进制协议序列化消息,而kafka-console-producer.sh发送的是纯文本字符串,两者格式不兼容,导致消费者反序列化失败。错误日志里的protocol id '104'其实是文本"hello"首字母'h'的ASCII值,直接印证了消费者在解析Quix专属协议时,收到的却是纯文本数据。

解决方案

两种可行方案,按需选择:

方案1:用Quix Streams发送兼容格式消息

如果要保留Quix的Timeseries DataFrame处理逻辑,就得用Quix生产者发送符合其格式的消息,替换控制台生产者。示例代码:

import quixstreams as qx
import pandas as pd
import time

client = qx.KafkaStreamingClient('127.0.0.1:9092')
topic_producer = client.get_topic_producer("quickstart-events")

# 创建流并发送带消息的DataFrame
stream = topic_producer.create_stream()
df = pd.DataFrame({
    "value": ["hello, world!"],
    "timestamp": [pd.Timestamp.now()]
})
stream.timeseries.publish(df)

print("消息发送成功")

方案2:修改消费者直接处理纯文本

如果要消费控制台生产者发的纯文本,绕过Quix的Timeseries解析逻辑,直接获取原始消息即可。修改后的消费者代码:

import quixstreams as qx

client = qx.KafkaStreamingClient('127.0.0.1:9092')
# 使用RawTopicConsumer获取原始字节数据
topic_consumer = client.get_raw_topic_consumer("quickstart-events", consumer_group=None)

def on_raw_message_received(msg: qx.RawMessage):
    # 字节转字符串(默认utf-8编码)
    text = msg.value.decode("utf-8")
    print(f"收到消息: {text}")

topic_consumer.on_message_received = on_raw_message_received

print("监听原始消息中,按CTRL-C退出")
qx.App.run()
验证方式
  • 方案1:运行Quix生产者代码后,原消费者会正常打印包含消息的DataFrame。
  • 方案2:用kafka-console-producer.sh发送消息后,修改后的消费者会直接输出纯文本内容。

内容的提问来源于stack exchange,提问作者kovac

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 01:45:32