Python Confluent Kafka消费无法获取消息键及内容解析异常问题
问题描述
使用Python的confluent-kafka库消费Kafka主题消息时,遇到两个问题:
- 获取到的消息值格式混乱,带有以H开头的异常前缀内容
- 无法正确获取消息中的键(原始消息包含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
相关产品推荐
相关产品推荐

