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

AWS Kinesis Firehose+Lambda:如何处理并发记录创建/更新问题

问题描述

我有一个作为Kinesis Firehose工作流组成部分的Lambda函数,关联到Kinesis数据流,负责处理来自DynamoDB的序列化记录,做反序列化和S3编解码操作,代码如下:

def lambda_handler(event, context):
    deserializer = TypeDeserializer()
    output = []
    try:
        for record in event["records"]:
            # Everything stored in S3 is base64 encoded, so we must first decode, deserialize, and
            # then encode the records again before we send them to S3
            decoded_payload = json.loads(base64.b64decode(record["data"]).decode())

            if decoded_payload["dynamodb"] and decoded_payload["dynamodb"]["NewImage"]:
                updated_record = decoded_payload["dynamodb"]["NewImage"]
                deserialized_record = {
                    k: deserializer.deserialize(v) for k, v in updated_record.items()
                }

                # Add newline after each record??? (otherwise Athena will only "see" the first?)
                encoded_record_with_line_break = (json.dumps(deserialized_record, cls=DecimalEncoder) + "\n").encode()
                output_record = {
                    "recordId": record["recordId"],
                    "result": "Ok",
                    "data": base64.b64encode(encoded_record_with_line_break).decode(),
                }
            else:
                print(
                    f"Dropping payload containing no updated DynamoDB image: {decoded_payload}"
                )
                output_record = {
                    "recordId": record["recordId"],
                    "result": "Dropped",
                    "data": record["data"],
                }

            output.append(output_record)

            print(f"Successfully processed {len(event['records'])} records.")
        except Exception as exc:
            print(f"Unhandled exception raised during data transformation: {exc}")
        finally:
            return {"records": output}

最初遇到的问题是多条记录同时更新时,Glue Crawler无法识别所有记录,因为Kinesis Firehose默认把JSON记录以行内形式发送,无分隔符,导致Athena只能查到每个S3对象里的第一条记录。按照建议在每条记录后加\n换行符后,新问题出现:CloudWatch日志显示重复记录,单条更新生成2个文件,Athena查询报HIVE_BAD_DATA格式错误。

问题分析与修复方案

1. 重复记录/多文件问题修复

重复记录是因为Kinesis Firehose的重试机制:如果Lambda超时、报错或Firehose未收到响应,会重新调用Lambda处理同一批次记录;单条更新生成多文件则可能是Firehose的缓冲阈值设置过小,触发了提前归档。

处理步骤:

  • 查看Lambda错误日志,确认是否有超时或异常导致重试;
  • 调整Firehose缓冲配置:增大Buffer Size(比如从5MB调至10MB)或延长Buffer Interval(比如从60秒调至300秒),减少小文件生成;
  • 在Lambda中添加幂等处理:用recordId做去重标识,避免同一记录被重复处理。

2. Athena HIVE_BAD_DATA错误修复

手动给单条记录加换行符可能导致Firehose写入S3时格式混乱,加上Decimal类型序列化不当也会生成无效JSON。正确的做法是让Firehose自动处理换行分隔,而非Lambda手动添加。

处理步骤:

  • 确保DecimalEncoder正确序列化DynamoDB的Decimal类型,避免JSON格式错误;
  • 关闭Lambda中手动添加换行符的逻辑,改为启用Firehose的**Record Format Conversion**功能:配置输出格式为JSON,并设置Row Format为JSON Lines,Firehose会自动在每条记录间添加换行符;
  • 若必须在Lambda中处理,确保输出的每个记录是独立合法的JSON字符串,且与Firehose的缓冲配置兼容。

修正后的Lambda代码示例

import json
import base64
from boto3.dynamodb.types import TypeDeserializer
from decimal import Decimal

class DecimalEncoder(json.JSONEncoder):
    def default(self, obj):
        if isinstance(obj, Decimal):
            return float(obj) if obj % 1 != 0 else int(obj)
        return super(DecimalEncoder, self).default(obj)

def lambda_handler(event, context):
    deserializer = TypeDeserializer()
    output = []
    processed_record_ids = set()
    
    try:
        for record in event["records"]:
            # 幂等处理:跳过已处理的recordId
            if record["recordId"] in processed_record_ids:
                output.append({
                    "recordId": record["recordId"],
                    "result": "Ok",
                    "data": record["data"]
                })
                continue
            processed_record_ids.add(record["recordId"])
            
            decoded_payload = json.loads(base64.b64decode(record["data"]).decode())
            
            if decoded_payload.get("dynamodb") and decoded_payload["dynamodb"].get("NewImage"):
                updated_record = decoded_payload["dynamodb"]["NewImage"]
                deserialized_record = {
                    k: deserializer.deserialize(v) for k, v in updated_record.items()
                }
                # 仅序列化JSON,不手动加换行符
                encoded_record = json.dumps(deserialized_record, cls=DecimalEncoder).encode()
                output_record = {
                    "recordId": record["recordId"],
                    "result": "Ok",
                    "data": base64.b64encode(encoded_record).decode()
                }
            else:
                print(f"Dropping payload with no updated DynamoDB image: {decoded_payload}")
                output_record = {
                    "recordId": record["recordId"],
                    "result": "Dropped",
                    "data": record["data"]
                }
            
            output.append(output_record)
        
        print(f"Successfully processed {len(output)} records.")
    except Exception as exc:
        print(f"Error during data transformation: {exc}")
        # 异常时返回失败记录,避免Firehose无限重试
        for record in event["records"]:
            if record["recordId"] not in processed_record_ids:
                output.append({
                    "recordId": record["recordId"],
                    "result": "ProcessingFailed",
                    "data": record["data"]
                })
    finally:
        return {"records": output}

额外配置要点

  1. Firehose格式转换配置:
    • 开启Record Format Conversion,选择Output format为JSON;
    • 在JSON configuration中设置Row format为JSON Lines,自动添加记录分隔符。
  2. Glue Crawler配置:
    • 确保分类器为JSON,且启用JSON Lines识别;
    • 设置定期运行或事件触发,及时更新表结构。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 09:00:42