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

使用DynamoDB Stream同步数据到Redshift遇全NULL值问题求助

问题:DynamoDB数据同步Redshift后全列为NULL的解决方法

问题背景

通过DynamoDB Stream、Lambda、Kinesis Delivery Firehose将DynamoDB数据同步至Redshift,所有流程无报错,但Redshift的profile表所有列均为NULL。相关信息如下:

DynamoDB数据结构

{ 
  "Id": { "S": "108" }, 
  "Name": { "S": "Smith" },  
  "Address": { 
    "M": {  
      "City":{ "S": "California" },
      "State":{ "S": "New York" },
      "Country":{ "S": "USA" } 
    } 
  },  
  "Email": { "S": "xyz2dffd@gmail.com" },  
  "Designation": { "S": "Sr.Dev" },  
  "PhoneNumber": { "S": "2124141241" } 
}

Redshift表Schema

PROFILE ( 
  Id varchar(255), 
  Name varchar(255),
  Address SUPER,   
  Email varchar(255),
  Designation varchar(255),
  PhoneNumber varchar(255) 
);

原Lambda代码

import json
import boto3

firehose = boto3.client('firehose')
deliveryStreamName = 'Delivery Stream Name'

def convertToFirehoseRecord(ddbRecord):
    newImage = ddbRecord['NewImage']
    address = newImage.get('Address', {}).get('M', {})
    firehoseRecord = {
        "Id": newImage.get('Id', {}).get('S'),
        "Name": newImage.get('Name', {}).get('S'),
        "Address": {
            "State": address.get('State', {}).get('S'),
            "Country": address.get('Country', {}).get('S'),
            "City": address.get('City', {}).get('S')
        },
        "Email": newImage.get('Email', {}).get('S'),
        "Designation": newImage.get('Designation', {}).get('S'),
        "PhoneNumber": newImage.get('PhoneNumber', {}).get('S')
    }
    return firehoseRecord

def lambda_handler(event, context):
    print(event)
    firehoseRecords = []
    for record in event['Records']:
        print(record)
        ddbRecord = record['dynamodb']
        print('DDB Record: ' + json.dumps(ddbRecord))
        
        firehoseRecord = convertToFirehoseRecord(ddbRecord)
        print('Firehose Record: ' + json.dumps(firehoseRecord))
        
        firehoseRecords.append({'Data': json.dumps(firehoseRecord)})
    print('list')
    print(firehoseRecords)
    result = firehose.put_record_batch(DeliveryStreamName=deliveryStreamName, Records=firehoseRecords)
    print(result)
    return 'processed {} records.'.format(len(event['Records']))

使用的COPY命令

COPY profile FROM 's3://my-bucket/<manifest>' 
CREDENTIALS 'aws_iam_role=arn:aws:iam::<aws-account-id>:role/<role-name>' 
MANIFEST json 'auto';

核心原因

Redshift的json 'auto'参数严格区分大小写,Lambda生成的JSON键为驼峰格式(如Id、Address),但Redshift表的列名实际是小写(id、address),导致字段匹配失败,所有列填充为NULL。

解决方法

1. 修正Lambda代码的JSON键为小写

将生成的Firehose记录的键改为全小写,与Redshift列名完全匹配:

import json
import boto3

firehose = boto3.client('firehose')
deliveryStreamName = 'Delivery Stream Name'

def convertToFirehoseRecord(ddbRecord):
    newImage = ddbRecord['NewImage']
    address_raw = newImage.get('Address', {}).get('M', {})
    # 键改为小写,匹配Redshift列名
    firehoseRecord = {
        "id": newImage.get('Id', {}).get('S'),
        "name": newImage.get('Name', {}).get('S'),
        "address": {
            "state": address_raw.get('State', {}).get('S'),
            "country": address_raw.get('Country', {}).get('S'),
            "city": address_raw.get('City', {}).get('S')
        },
        "email": newImage.get('Email', {}).get('S'),
        "designation": newImage.get('Designation', {}).get('S'),
        "phonenumber": newImage.get('PhoneNumber', {}).get('S')
    }
    return firehoseRecord

def lambda_handler(event, context):
    firehoseRecords = []
    for record in event['Records']:
        ddbRecord = record['dynamodb']
        # 过滤删除/旧数据记录,只处理有NewImage的新增/更新
        if 'NewImage' not in ddbRecord:
            continue
            
        firehoseRecord = convertToFirehoseRecord(ddbRecord)
        # 每条记录末尾添加换行符,确保Redshift能正确解析单行JSON
        firehoseRecords.append({'Data': json.dumps(firehoseRecord) + '\n'})
    
    if firehoseRecords:
        result = firehose.put_record_batch(DeliveryStreamName=deliveryStreamName, Records=firehoseRecords)
        # 检查批量写入是否有失败记录
        failed_count = result.get('FailedPutCount', 0)
        if failed_count > 0:
            raise Exception(f"Failed to put {failed_count} records to Firehose: {result['RequestResponses']}")
    
    return f'processed {len(firehoseRecords)} valid records.'

2. 验证Firehose输出的JSON格式

确保Firehose写入S3的文件中,每条JSON记录单独一行,格式如下:

{"id":"108","name":"Smith","address":{"state":"New York","country":"USA","city":"California"},"email":"xyz2dffd@gmail.com","designation":"Sr.Dev","phonenumber":"2124141241"}

3. 重新执行COPY命令

保持原COPY命令不变,重新同步数据即可。

代码优化建议

Python版本优化

  • 使用boto3.dynamodb.types.TypeDeserializer自动解析DynamoDB的原生格式,避免手动嵌套get:
from boto3.dynamodb.types import TypeDeserializer

def convertToFirehoseRecord(ddbRecord):
    deserializer = TypeDeserializer()
    # 自动解析DDB的NewImage为普通Python字典
    parsed_item = {k: deserializer.deserialize(v) for k, v in ddbRecord['NewImage'].items()}
    # 转换键为小写,处理Address嵌套结构
    return {
        "id": parsed_item.get('Id'),
        "name": parsed_item.get('Name'),
        "address": {
            "state": parsed_item.get('Address', {}).get('State'),
            "country": parsed_item.get('Address', {}).get('Country'),
            "city": parsed_item.get('Address', {}).get('City')
        },
        "email": parsed_item.get('Email'),
        "designation": parsed_item.get('Designation'),
        "phonenumber": parsed_item.get('PhoneNumber')
    }
  • 添加异常捕获,处理DynamoDB Stream的不同事件类型(如删除事件无NewImage)
  • 批量写入Firehose时,控制单次批量的记录数(Firehose批量上限为500条)

Node.js版本参考

const AWS = require('aws-sdk');
const firehose = new AWS.Firehose();
const deliveryStreamName = 'Delivery Stream Name';

// 解析DynamoDB原生格式为普通对象
function deserializeDdbItem(item) {
  return Object.entries(item).reduce((acc, [key, value]) => {
    const type = Object.keys(value)[0];
    acc[key] = value[type];
    return acc;
  }, {});
}

function convertToFirehoseRecord(ddbRecord) {
  const parsedItem = deserializeDdbItem(ddbRecord.NewImage);
  return {
    id: parsedItem.Id,
    name: parsedItem.Name,
    address: {
      state: parsedItem.Address?.State,
      country: parsedItem.Address?.Country,
      city: parsedItem.Address?.City
    },
    email: parsedItem.Email,
    designation: parsedItem.Designation,
    phonenumber: parsedItem.PhoneNumber
  };
}

exports.handler = async (event) => {
  const firehoseRecords = [];
  for (const record of event.Records) {
    const ddbRecord = record.dynamodb;
    if (!ddbRecord.NewImage) continue;
    
    const firehoseRecord = convertToFirehoseRecord(ddbRecord);
    firehoseRecords.push({
      Data: JSON.stringify(firehoseRecord) + '\n'
    });
  }
  
  if (firehoseRecords.length > 0) {
    const result = await firehose.putRecordBatch({
      DeliveryStreamName: deliveryStreamName,
      Records: firehoseRecords
    }).promise();
    
    if (result.FailedPutCount > 0) {
      throw new Error(`Failed to put ${result.FailedPutCount} records: ${JSON.stringify(result.RequestResponses)}`);
    }
  }
  
  return `processed ${firehoseRecords.length} valid records.`;
};

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 06:10:05