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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 23:20:40