使用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
相关产品推荐
相关产品推荐

