使用AWS Lambda解码Kinesis数据流时遭遇UTF-8解码错误
AWS Lambda解码Kinesis数据流时的UTF-8解码错误解决方法
问题描述
我在AWS Lambda中尝试解码AWS Kinesis数据流的数据时,持续收到错误:'utf-8' codec can't decode bytes in position 0-2: invalid continuation byte,相关代码如下:
x = b'\xf3\x89\x9a\xc2\n$dad568a5-6305-481c-b6f1-f8338cc127df\n$3d57f33a-d681-467b-bb82-89c0d77e2621\n$3ade7757-3df4-41ec-bdc8-52a27449c420\n$a0a59a4e-02f5-462d-8c3e-50030145cf17\x1a\x83\x01\x08\x00\x1a\x7f{ "window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2004,"bonus_stake": 2.76,"bonus_win": 4}\x1a\x86\x01\x08\x01\x1a\x81\x01{"window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2304,"bonus_stake": 2.2,"bonus_win": 2.21}\x1a\x87\x01\x08\x02\x1a\x82\x01{"window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2290,"bonus_stake": 11.1,"bonus_win": 38.7}\x1a\x86\x01\x08\x03\x1a\x81\x01{"window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2192,"bonus_stake": 1.32,"bonus_win": 0.6}\x10\xa6\x1a\tB\xa5\x9b\x14\xa5?\xad\xcd\x8b\xe8^\xcb' s = x.decode() print(s)
解决方案
问题根源
你的字节数据开头包含非UTF-8兼容的二进制前缀(\xf3\x89\x9a\xc2),还有中间的控制字符(如\x1a、\x08等),直接用默认UTF-8解码会失败。这些前缀和控制字符通常是Kinesis数据的序列化协议头(比如Protobuf、Avro或自定义二进制格式),不是纯UTF-8文本。
具体解决步骤
- 确认上游序列化格式:先明确发送到Kinesis的数据使用的序列化方式(如Protobuf、Avro),如果是Protobuf需要对应的
.proto文件解析,Avro则需要Schema。 - 提取有效JSON内容:从你的字节数据来看,包含完整的JSON片段,可以定位有效JSON起始位置,跳过二进制前缀:
x = b'\xf3\x89\x9a\xc2\n$dad568a5-6305-481c-b6f1-f8338cc127df\n$3d57f33a-d681-467b-bb82-89c0d77e2621\n$3ade7757-3df4-41ec-bdc8-52a27449c420\n$a0a59a4e-02f5-462d-8c3e-50030145cf17\x1a\x83\x01\x08\x00\x1a\x7f{ "window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2004,"bonus_stake": 2.76,"bonus_win": 4}\x1a\x86\x01\x08\x01\x1a\x81\x01{"window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2304,"bonus_stake": 2.2,"bonus_win": 2.21}\x1a\x87\x01\x08\x02\x1a\x82\x01{"window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2290,"bonus_stake": 11.1,"bonus_win": 38.7}\x1a\x86\x01\x08\x03\x1a\x81\x01{"window_start": "2022-12-30 13:25:00","window_end": "2022-12-30 13:35:00","player_id": 2192,"bonus_stake": 1.32,"bonus_win": 0.6}\x10\xa6\x1a\tB\xa5\x9b\x14\xa5?\xad\xcd\x8b\xe8^\xcb' # 定位第一个JSON的起始位置 json_start = x.find(b'{') if json_start != -1: # 替换HTML实体为双引号,提取有效JSON字节 json_bytes = x[json_start:].replace(b'"', b'"') # 忽略末尾无效字节完成解码 s = json_bytes.decode('utf-8', errors='ignore') print(s) - 使用对应序列化库解析:如果上游用了Protobuf,安装
protobuf库并导入消息类解析;如果是Avro,用fastavro或avro库加载Schema后解析。 - 处理Kinesis的Base64编码:Lambda触发器接收的Kinesis记录通常是Base64编码的,需先解码再处理:
import base64 def lambda_handler(event, context): for record in event['Records']: # 解码Base64格式的Kinesis数据 payload = base64.b64decode(record['kinesis']['data']) # 后续按序列化格式解析payload # ...
内容的提问来源于stack exchange,提问作者Pregz
相关产品推荐
相关产品推荐

