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}
额外配置要点
- Firehose格式转换配置:
- 开启
Record Format Conversion,选择Output format为JSON; - 在
JSON configuration中设置Row format为JSON Lines,自动添加记录分隔符。
- 开启
- Glue Crawler配置:
- 确保分类器为
JSON,且启用JSON Lines识别; - 设置定期运行或事件触发,及时更新表结构。
- 确保分类器为
内容的提问来源于stack exchange,提问作者tprebenda
相关产品推荐
相关产品推荐

