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

Kinesis Firehose Lambda数据转换响应JSON序列化报错排查

Kinesis Data Firehose Lambda转换函数报错解决方案

问题根因

  • base64.b64encode() 方法返回的是bytes类型对象,Lambda返回响应时需要序列化为JSON格式,JSON不支持bytes类型,直接返回会触发序列化报错。
  • output变量定义为全局变量,Lambda运行容器复用时会累计之前调用的处理记录,会导致数据重复、内存溢出等问题。
  • 直接使用str()转换Python字典为字符串,得到的是Python原生字典格式(键值用单引号包裹),不属于标准JSON格式,后续Athena读取时也会出现解析错误。

修正后完整代码

#This function is created to Transform the data from Kinesis Data Firehose -> S3 Bucket
#It converts single line json to multi line json as expected by AWS Athena best practice.
#It also removes special characters from json keys (column name in Athena) as Athena expects column names without special characters

import json
import base64
import string
from typing import Optional, Iterable, Union

delete_dict = {sp_character: '' for sp_character in string.punctuation}
PUNCT_TABLE = str.maketrans(delete_dict)

def lambda_handler(event, context):
    output = []
    for record in event['records']:
        payload = base64.b64decode(record['data']).decode('utf-8')
        
        remove_special_char = json.loads(payload, object_pairs_hook=clean_keys)
        # 用json.dumps生成标准JSON字符串,换行后再做base64编码
        row_w_newline = json.dumps(remove_special_char) + "\n"
        row_w_newline = base64.b64encode(row_w_newline.encode('utf-8')).decode('utf-8')
        
        output_record = {
            'recordId': record['recordId'],
            'result': 'Ok',
            'data': row_w_newline
        }
        output.append(output_record)

    print('Processed {} records.'.format(len(event['records'])))
    
    return {'records': output}
    
def strip_punctuation(s: str,
                      exclude_chars: Optional[Union[str, Iterable]] = None) -> str:
    """
    Remove punctuation and spaces from a string.

    If `exclude_chars` is passed, certain characters will not be removed
    from the string.

    """
    punct_table = PUNCT_TABLE.copy()
    if exclude_chars:
        for char in exclude_chars:
            punct_table.pop(ord(char), None)

    # Next, remove the desired punctuation from the string
    return s.translate(punct_table)

def clean_keys(o):
    return {strip_punctuation(k): v for k, v in o}

核心修正说明

  • 将output变量移入lambda_handler函数内部,每次函数调用重新初始化,避免容器复用导致的数据累计问题
  • 使用json.dumps()替代str()处理字典,生成符合规范的JSON字符串,适配Athena解析要求
  • base64编码完成后调用decode('utf-8')将bytes类型转为字符串,解决JSON序列化报错问题
  • 移除了未使用的boto3导入,减少函数初始化开销

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 05:54:01