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

求助:C#版AWS Lambda同步DynamoDB删除数据至S3遇格式与权限错误

问题解决:Kinesis Firehose调用Lambda时的格式错误与AccessDenied问题

问题场景与报错信息

已搭建AWS Lambda(C#编写)、Kinesis Firehose和S3桶,用于将DynamoDB删除数据保存至S3。Lambda核心代码如下:

var client = new AmazonKinesisFirehoseClient();
try
{
    context.Logger.LogInformation($"Write to Kinesis Firehose: {list.Count}");
    var request = new PutRecordBatchRequest
    {
        DeliveryStreamName = _kinesisStream,
        Records = new List<Amazon.KinesisFirehose.Model.Record> ()
    };
    foreach (var item in list)
    {
        var stringWrite = new StringWriter();
        
        string json = JsonConvert.SerializeObject(item, new JsonSerializerSettings { NullValueHandling = NullValueHandling.Ignore });
        byte[] byteArray = UTF8Encoding.UTF8.GetBytes(ToLiteral(json));
        var record = new Amazon.KinesisFirehose.Model.Record
        {
            Data = new MemoryStream(byteArray)
        };
        request.Records.Add(record);
    }
    if (request.Records.Count > 0)
    {
        var response = await client.PutRecordBatchAsync(request);
        Console.WriteLine($"FailedPutCount: {response.FailedPutCount} status: {response.HttpStatusCode}");
    }
}
catch (Exception e)
{
    Console.WriteLine(e);
}

其中list为待处理的对象列表。当前遇到两个核心问题:

  1. Kinesis Firehose日志报错:

"message": "Check your function and make sure the output is in required format. In addition to that, make sure the processed records contain valid result status of Dropped, Ok, or ProcessingFailed",
"errorCode": "Lambda.FunctionError"

  1. S3桶出现AccessDenied错误,下载错误文件后可见:

"attemptsMade":4,"arrivalTimestamp":1661897166462,"errorCode":"Lambda.FunctionError","errorMessage":"Check your function and make sure the output is in required format. In addition to that, make sure the processed records contain valid result status of Dropped, Ok, or ProcessingFailed","attemptEndingTimestamp":1661897241573,"rawData":"XXXXXXXXX"

rawData解码后为写入Firehose的原始JSON字符串。

核心问题与修复方案

1. 错误根源:Lambda角色定位混淆

如果你的Lambda是作为Kinesis Firehose的数据转换Lambda(即Firehose配置了Lambda处理流入数据),Lambda的职责是接收Firehose传入的记录,处理后返回符合指定格式的响应,而非主动调用PutRecordBatchAsync写入Firehose。你的代码逻辑混淆了Lambda的角色,导致输出不符合Firehose要求的格式,触发Lambda.FunctionError。

2. 具体修复步骤

步骤1:调整Lambda输入输出格式

Firehose触发的Lambda必须遵循严格的输入输出规范:

  • 输入:接收KinesisFirehoseEvent对象,包含待处理的记录集合
  • 输出:必须返回包含records数组的对象,每个元素需指定recordId和result(可选data),result仅允许Ok、Dropped或ProcessingFailed三个值

修正后的C#代码示例:

using Amazon.Lambda.Core;
using Amazon.Lambda.KinesisFirehoseEvents;
using Newtonsoft.Json;
using System.Collections.Generic;
using System.Text;

[assembly: LambdaSerializer(typeof(Amazon.Lambda.Serialization.SystemTextJson.DefaultLambdaJsonSerializer))]

namespace FirehoseTransformLambda
{
    public class Function
    {
        public KinesisFirehoseResponse FunctionHandler(KinesisFirehoseEvent input, ILambdaContext context)
        {
            var response = new KinesisFirehoseResponse
            {
                Records = new List<KinesisFirehoseResponse.RecordResult>()
            };

            foreach (var record in input.Records)
            {
                try
                {
                    // 解码Firehose传入的Base64格式原始数据
                    string rawData = Encoding.UTF8.GetString(record.Data);
                    // 解析DynamoDB删除事件数据
                    var dynamoItem = JsonConvert.DeserializeObject<YourDynamoModel>(rawData);
                    
                    // 处理数据(此处可根据需求修改字段、过滤内容等)
                    string processedJson = JsonConvert.SerializeObject(dynamoItem, new JsonSerializerSettings { NullValueHandling = NullValueHandling.Ignore });
                    byte[] processedBytes = Encoding.UTF8.GetBytes(processedJson);
                    string encodedData = System.Convert.ToBase64String(processedBytes);

                    // 返回处理成功的结果
                    response.Records.Add(new KinesisFirehoseResponse.RecordResult
                    {
                        RecordId = record.RecordId,
                        Result = "Ok",
                        Data = encodedData
                    });
                }
                catch (System.Exception ex)
                {
                    context.Logger.LogError($"处理记录失败: {record.RecordId}, 错误信息: {ex.Message}");
                    // 标记处理失败,Firehose会自动重试或存入错误桶
                    response.Records.Add(new KinesisFirehoseResponse.RecordResult
                    {
                        RecordId = record.RecordId,
                        Result = "ProcessingFailed"
                    });
                }
            }

            return response;
        }
    }

    // 自定义DynamoDB数据模型,匹配删除事件的字段结构
    public class YourDynamoModel
    {
        public string Id { get; set; }
        // 新增其他对应字段...
    }
}

步骤2:修正IAM权限配置

S3的AccessDenied错误通常由权限不足导致:

  • 确保Firehose的IAM角色拥有s3:PutObject权限,且目标S3桶的桶策略允许该角色写入
  • 若Lambda需要访问DynamoDB等其他资源,确保Lambda的IAM角色拥有对应访问权限

步骤3:验证Firehose配置

确认Firehose的"数据转换"配置已正确关联目标Lambda,且未设置冲突的格式转换规则。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 15:48:11