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

Python Confluent Kafka消费无法获取消息键及内容解析异常问题

问题描述

使用Python的confluent-kafka库消费Kafka主题消息时,遇到两个问题:

  1. 获取到的消息值格式混乱,带有以H开头的异常前缀内容
  2. 无法正确获取消息中的键(原始消息包含id、operation等结构化键值对)

当前代码

from confluent_kafka import Consumer, KafkaError
import uuid

# Kafka broker details
broker = "something"
topic = "something"
group = str(uuid.uuid4())

# Kafka consumer configuration
conf = {
    "bootstrap.servers": broker,
    "group.id": group,
    "auto.offset.reset": "earliest"  
}

consumer = Consumer(conf)

# Subscribe to the topic
consumer.subscribe([topic])

try:
    while True:
        msg = consumer.poll(1.0)  # Poll for new messages
        
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                # End of partition event
                print("Reached end of partition")
            else:
                print("Error: {}".format(msg.error().str()))
        else:
            # Decode the received message value as UTF-8
            received_message = msg.value().decode("utf-8", errors="ignore")

            # Print the received message
            print("Received message:", received_message)
except KeyboardInterrupt:
    pass
finally:
    consumer.close()

当前输出

Received message: 
H2ff4c4d5-74ee-4708-b553-aa68ea5e675bHdfb821d1-db91-44b8-a080-8f689a80bb32H2ff4c4d5-74ee-4708-b553-aa68ea5e675bCREATEIntegrationCATALOGNodeܥ[{"id":"2ff4c4d5-74ee-4708-b553-aa68ea5e675b","name":"Integration","classType":"EntityClass","graphType":"Node","objectType":"Integration","enabled":true,"tenantDefault":true.....

原始Kafka消息(含完整键值)

"{\"id\":\"100dea7f-bc9d-4775-b660-17b99fd1439a\",\"tenant_id\":\"d257587b-d53f-498a-9256-4d1851b847d8\",\"callback_id\":{\"string\":\"100dea7f-bc9d-4775-b660-17b99fd1439a\"},\"operation\":{\"string\":\"CREATE\"},\"type\":{\"string\":\"Genericentity\"},\"base_type\":{\"string\":\"CATALOG\"},\"graph_type\":{\"string\":\"Node\"},\"json_payload\":{\"string\":\"[{\\\"id\\\":\\\"100dea7f-bc9d-4775-b660-17b99fd1439a\\\",\\\"name\\\":\\\"LSH - Safety Stock\\\",\\\"classType\\\":\\\"EntityClass\\\",\\\"graphType\\\":\\\"Node\\\",\\\"objectType\\\":\\\"LSHSafetyStock\\\",\\\"enabled\\\":true,\\\"tenantDefault\\\":false,\\\"catalogType\\\":\\\"genericentity\\\",\\\"type\\\":\\\"Genericentity\\\",\\\"cdmFields\\\":[{\\\"id\\\":\\\"2e1f4381-cf61-49a5-934b-6b47a93d1771\\\",\\\"name\\\":\\\"baseTypeKey\\\",\\\"displayLabel\\\":\\\"Base Type Key\\\",\\\"type\\\":\\\"STRING\\\",\\\"required\\\":false,\\\"size\\\":20,\\\"unit\\\":\\\"Default\\\",\\\"value\\\":\\\"\\\",\\\"regex\\\":null,\\\"unique\\\":false,\\\"defaultValue\\\":null,\\\"group\\\":\\\"System\\\",\\\"autoGenerated\\\":false,\\\"format\\\":\\\"\\\",\\\"tooltip\\\":\\\"\\\",\\\"searchable\\\":false,\\\"filterable\\\":false,\\\"editable\\\":false,\\\"entityName\\\":null,\\\"entityKey\\\":null,\\\"entityValue\\\":null,\\\"displayable\\\":false,\\\"instanceUserEditable\\\":false,\\\"instanceUserCreatable\\\":false,\\\"apiFlag\\\":false,\\\"apiUrl\\\":\\\"\\\",\\\"apiResponseSelectKey\\\":\\\"\\\",\\\"apiResponseSelectValue\\\":\\\"\\\",\\\"order\\\":-1,\\\"tenantDefault\\\":false,\\\"isDisplayableOnSummary\\\":false,\\\"isDisplayableOnDetails\\\":false,\\\"isDisplayableOnCatalog\\\":false,\\\"isDisplayableOnList\\\":false,\\\"dataClassification\\\":\\\"DEFAULT\\\"},{\\\"id\\\":\\\"cf7a08cb-5f7f-472b-a9eb-516ca864a188\\\",\\\"name\\\":\\\"entityType\\\",\\\"displayLabel\\\":\\\"Entity Type\\\",\\\"type\\\":\\\"STRING\\\",\\\"required\\\":false,\\\"size\\\":20,\\\"unit\\\":\\\"Default\\\",\\\"value\\\":\\\"\\\",\\\"regex\\\":null,\\\"unique\\\":false,\\\"defaultValue\\\":null,\\\"group\\\":\\\"System\\\",\\\"autoGenerated\\\":false,\\\"format\\\":\\\"\\\",\\\"tooltip\\\":\\\"\\\",\\\"searchable\\\":false,\\\"filterable\\\":false,\\\"editable\\\":false,\\\"entityName\\\":null,\\\"entityKey\\\":null,\\\"entityValue\\\":null,\\\"displayable\\\":false,\\\"instanceUserE
解决方案

1. 处理消息值的异常前缀与序列化格式

从输出看,消息值并非纯UTF-8字符串,而是二进制序列化格式(比如Avro、Protobuf或自定义二进制协议),直接UTF-8解码会导致乱码和前缀问题。

临时解决:提取消息中的JSON部分

观察到输出末尾存在合法JSON结构,可通过字符串匹配提取:

import json

raw_value = msg.value()
# 找到第一个{的位置,提取后续JSON内容
json_start = raw_value.find(b'{')
if json_start != -1:
    json_part = raw_value[json_start:]
    try:
        parsed_msg = json.loads(json_part.decode('utf-8'))
        print("解析后的消息值:", parsed_msg)
    except json.JSONDecodeError as e:
        print("JSON解析失败:", e)

长期解决:匹配生产者序列化协议

必须明确生产者使用的序列化方式,才能完全正确解析:

  • 如果是Avro:使用confluent-kafka[avro]库的AvroConsumer,并配置Schema Registry地址
  • 如果是Protobuf:使用对应Protobuf库解析二进制数据
  • 如果是自定义格式:按照生产者的编码规则编写解析逻辑

2. 获取消息的键

confluent-kafka的Message对象提供key()方法获取消息键,需根据键的序列化方式解码:

# 在处理消息的else分支中添加
msg_key = msg.key()
if msg_key is not None:
    # 如果键是UTF-8字符串
    decoded_key = msg_key.decode('utf-8')
    print("消息键:", decoded_key)
    # 如果键是JSON格式
    # parsed_key = json.loads(msg_key.decode('utf-8'))
    # print("解析后的消息键:", parsed_key)

3. 完整优化后的代码

from confluent_kafka import Consumer, KafkaError
import uuid
import json

# Kafka broker details
broker = "something"
topic = "something"
group = str(uuid.uuid4())

# Kafka consumer configuration
conf = {
    "bootstrap.servers": broker,
    "group.id": group,
    "auto.offset.reset": "earliest"  
}

consumer = Consumer(conf)
consumer.subscribe([topic])

try:
    while True:
        msg = consumer.poll(1.0)
        
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                print("已到达分区末尾")
            else:
                print("错误: {}".format(msg.error().str()))
        else:
            # 处理消息键
            msg_key = msg.key()
            if msg_key:
                try:
                    decoded_key = msg_key.decode('utf-8')
                    print("消息键:", decoded_key)
                except UnicodeDecodeError:
                    print("消息键无法UTF-8解码,可能是二进制格式")
            
            # 处理消息值
            raw_value = msg.value()
            if raw_value:
                # 尝试提取JSON部分
                json_start = raw_value.find(b'{')
                if json_start != -1:
                    try:
                        parsed_value = json.loads(raw_value[json_start:].decode('utf-8'))
                        print("解析后的消息值:", parsed_value)
                    except json.JSONDecodeError as e:
                        print("JSON解析失败:", e)
                else:
                    print("未找到有效JSON结构,原始值:", raw_value.decode('utf-8', errors='ignore'))
except KeyboardInterrupt:
    pass
finally:
    consumer.close()

关键说明

  • 核心问题是消息序列化格式不匹配:生产者用二进制序列化发送,消费者用UTF-8直接解码导致乱码
  • 临时提取JSON只是权宜之计,长期必须匹配生产者的序列化协议
  • 消息键的获取依赖msg.key()方法,需根据实际编码方式解码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 20:23:12