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

如何将格式错误的字节数据转换为字典?(Kafka场景)

解决Kafka消费者解析包含datetime的字节字典报错问题

咱们先搞清楚为啥ast.literal_eval()会报错:这个方法只能解析Python的字面量类型(比如字典、列表、字符串、数字这些),但你的消息字节串里包含了datetime.datetime(2018, 5, 25, 13, 35, 11)——这是一个类实例的构造调用,不属于字面量范畴,所以解析失败是必然的。

下面给你几种可行的解决方案,按安全和通用性排序:

方案一:用JSON序列化(推荐)

JSON是跨语言、安全的序列化格式,唯一需要处理的是datetime对象不能直接被JSON序列化,咱们自定义一个编码器就行。

修改生产者代码

把原来直接转bytes的逻辑改成JSON序列化:

import json
from datetime import datetime

# 自定义JSON编码器,处理datetime对象
def serialize_datetime(obj):
    if isinstance(obj, datetime):
        return obj.isoformat()  # 转成ISO标准的字符串格式,方便后续还原

# 替换原有的send逻辑
try:
    self.idSelf += 1
    self.tweet["tweet"] = tweetText
    self.tweet["id"] = id
    self.tweet["sequence"] = self.idSelf
    self.tweet["created_at"] = created_at
    # 先序列化为JSON字符串,再转成bytes发送
    tweet_json_str = json.dumps(self.tweet, default=serialize_datetime)
    future = producer.send('bubble-test', tweet_json_str.encode('utf8'))
except Exception as e:
    print(e)
finally:
    producer.flush()

修改消费者代码

解析JSON字符串,再把created_at还原为datetime对象:

import json
from datetime import datetime
import sys

consumer = KafkaConsumer('bubble-test', bootstrap_servers='localhost:9092', auto_offset_reset='earliest')
for message in consumer:
    # 先把bytes转成字符串
    tweet_json_str = message.value.decode('utf8')
    # 解析为字典
    tweet = json.loads(tweet_json_str)
    # 把ISO字符串转回datetime对象
    tweet['created_at'] = datetime.fromisoformat(tweet['created_at'])
    print(type(tweet))
    print(tweet)
    sys.exit()

方案二:自定义AST节点解析(不推荐,太繁琐)

如果一定要保留原有的生产者序列化方式,你可以手动解析AST节点,处理datetime的构造调用,但这个方法需要自己处理各种节点类型,维护成本高:

import ast
from datetime import datetime
import sys

def eval_ast_node(node):
    # 处理datetime构造调用
    if isinstance(node, ast.Call) and hasattr(node.func, 'id') and node.func.id == 'datetime.datetime':
        args = [eval_ast_node(arg) for arg in node.args]
        return datetime(*args)
    # 处理字典
    elif isinstance(node, ast.Dict):
        return {eval_ast_node(k): eval_ast_node(v) for k, v in zip(node.keys, node.values)}
    # 处理常量(字符串、数字等)
    elif isinstance(node, ast.Constant):
        return node.value
    # 其他类型按需扩展
    else:
        raise ValueError(f"无法处理的节点类型: {type(node)}")

consumer = KafkaConsumer('bubble-test', bootstrap_servers='localhost:9092', auto_offset_reset='earliest')
for message in consumer:
    tweet_str = message.value.decode('utf8')
    # 解析AST树
    tree = ast.parse(tweet_str, mode='eval')
    tweet = eval_ast_node(tree.body)
    print(type(tweet))
    print(tweet)
    sys.exit()

方案三:用Pickle序列化(不推荐,有安全风险)

Pickle可以直接序列化Python对象,但只能在完全信任的生产者和消费者之间使用——因为Pickle可以执行任意代码,如果消息被篡改,会导致严重的安全漏洞。

生产者代码修改

import pickle

# 替换原有的send逻辑
try:
    self.idSelf += 1
    self.tweet["tweet"] = tweetText
    self.tweet["id"] = id
    self.tweet["sequence"] = self.idSelf
    self.tweet["created_at"] = created_at
    # 用pickle序列化对象
    tweet_bytes = pickle.dumps(self.tweet)
    future = producer.send('bubble-test', tweet_bytes)
except Exception as e:
    print(e)
finally:
    producer.flush()

消费者代码修改

import pickle
import sys

consumer = KafkaConsumer('bubble-test', bootstrap_servers='localhost:9092', auto_offset_reset='earliest')
for message in consumer:
    tweet = pickle.loads(message.value)
    print(type(tweet))
    print(tweet)
    sys.exit()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:01:07