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

AWS Glue迁移DynamoDB:如何保留空列而非转为Null值

DynamoDB跨表迁移时避免缺失字段被设为NULL的解决方法

出现这个问题的核心原因是所用迁移工具默认给不存在的字段填充了NULL类型值,要让目标表和源表完全一致,关键是只复制源表实际存在的字段,不额外添加任何字段。以下是几种可行的实现方式:

一、自定义Lambda函数迁移(最直接可控)

直接通过代码遍历源表数据,原样写入目标表,完全保留源表Item的结构:

import boto3

# 初始化DynamoDB资源
dynamodb = boto3.resource('dynamodb')
source_table = dynamodb.Table('你的源表名称')
target_table = dynamodb.Table('你的目标表名称')

def migrate_data():
    # 分页扫描源表,避免数据量大时超时
    paginator = source_table.get_paginator('scan')
    for page in paginator.paginate():
        # 批量写入目标表
        with target_table.batch_writer() as batch:
            for item in page['Items']:
                # 直接写入源表的原始Item,不做任何字段补充
                batch.put_item(Item=item)

if __name__ == '__main__':
    migrate_data()

注意:如果源表数据量极大,需要调整Lambda的超时时间(最长15分钟)和内存配置,或者分段处理数据。

二、AWS Glue迁移(自定义ETL逻辑过滤NULL)

如果用Glue做迁移,默认会把缺失字段转为NULL,需要在ETL脚本中添加过滤逻辑,移除值为NULL的字段:

from awsglue.context import GlueContext
from pyspark.context import SparkContext

sc = SparkContext()
glueContext = GlueContext(sc)

# 读取源表数据
source_dyf = glueContext.create_dynamic_frame.from_options(
    connection_type="dynamodb",
    connection_options={"dynamodb.input.tableName": "你的源表名称"}
)

# 定义转换函数:移除所有值为NULL的字段
def clean_null_fields(record):
    return {key: val for key, val in record.items() if val is not None}

# 应用转换逻辑
cleaned_dyf = source_dyf.map(clean_null_fields)

# 写入目标表
glueContext.write_dynamic_frame.from_options(
    frame=cleaned_dyf,
    connection_type="dynamodb",
    connection_options={"dynamodb.output.tableName": "你的目标表名称"}
)

三、S3导出+预处理+导入(适合超大规模数据)

如果源表数据量超大,先导出到S3再处理:

  1. 把源表数据导出到S3,选择JSON格式(避免CSV的NULL填充问题)
  2. 用Lambda或EMR处理S3中的JSON文件,过滤掉不存在的字段
  3. 将处理后的文件导入目标表

Lambda处理S3文件的示例代码:

import boto3
import json

s3 = boto3.client('s3')
dynamodb = boto3.resource('dynamodb')
target_table = dynamodb.Table('你的目标表名称')

def process_s3_file(event, context):
    for record in event['Records']:
        bucket = record['s3']['bucket']['name']
        file_key = record['s3']['object']['key']
        
        # 读取S3中的JSON文件
        s3_response = s3.get_object(Bucket=bucket, Key=file_key)
        items = json.loads(s3_response['Body'].read().decode('utf-8'))
        
        # 批量写入目标表,过滤NULL字段
        with target_table.batch_writer() as batch:
            for item in items:
                cleaned_item = {k: v for k, v in item.items() if v is not None}
                batch.put_item(Item=cleaned_item)

内容的提问来源于stack exchange,提问作者Shameem Banu Shaik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:02:47