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

基于CDK v2:无Glue实现DynamoDB流转自定义命名的S3 Parquet

解决方案:DynamoDB动态Schema数据转Parquet到S3(CDK v2实现)

架构选型

直接采用DynamoDB Stream + Lambda组合,放弃Kinesis Firehose+Glue方案——因为Glue Catalog要求固定Schema,无法适配你每个用户不同的data字段结构。若后续流量暴涨,可在DynamoDB与Lambda之间加一层Kinesis Data Streams做流量削峰,当前需求用DynamoDB原生Stream足够。

Lambda核心逻辑

Lambda需完成数据提取、分组、动态Schema生成、缓冲写入四个核心步骤:

  1. 提取有效数据:从DynamoDB Stream事件中解析出data字段及关联的用户ID(假设用户ID存在于DynamoDB条目的顶层属性,如userId)
  2. 按用户ID分组:将同用户的data数据聚合,避免跨用户混合Schema
  3. 动态生成Parquet Schema:用PyArrow或FastParquet库,从用户的data字典自动推断Schema(PyArrow的pa.from_pylist()会自动识别字段类型),生成带嵌入式Schema的Parquet文件
  4. 缓冲与写入S3:
    • 利用Lambda批处理触发配置:设置批量大小(如500条)和批处理窗口(如60秒),让Lambda攒够数据再调用
    • 内部缓冲:每个用户组累积到一定阈值(比如200条或10MB)再写入S3,避免生成大量小文件
    • 文件名构造:按user-data/{user_id}/YYYY-MM-DD/{uuid}.parquet格式命名,用用户ID做前缀、日期分区,uuid保证文件名唯一

Lambda代码片段(Python示例)

import boto3
import pyarrow as pa
import pyarrow.parquet as pq
import uuid
from datetime import datetime

s3 = boto3.client('s3')
BUCKET_NAME = 'your-target-bucket'

def lambda_handler(event, context):
    # 按用户ID分组数据
    user_data_groups = {}
    for record in event['Records']:
        # 解析DynamoDB新条目数据
        new_image = record['dynamodb']['NewImage']
        user_id = new_image['userId']['S']  # 假设用户ID为字符串类型
        data = new_image['data']['M']  # 解析JSON格式的data字段
        
        # 转换为Python字典适配PyArrow
        parsed_data = {}
        for k, v in data.items():
            if v.get('S'):
                parsed_data[k] = v['S']
            elif v.get('N'):
                parsed_data[k] = float(v['N']) if '.' in v['N'] else int(v['N'])
            elif v.get('BOOL'):
                parsed_data[k] = v['BOOL']
        
        # 分组存储
        if user_id not in user_data_groups:
            user_data_groups[user_id] = []
        user_data_groups[user_id].append(parsed_data)
    
    # 处理每个用户组的写入
    for user_id, data_list in user_data_groups.items():
        # 阈值判断:累积超过200条再写入,可根据需求调整
        if len(data_list) < 200:
            # 可选:若需持久化缓冲,可写入临时DynamoDB表,后续定时处理
            continue
        
        # 动态生成Schema并写入Parquet临时文件
        table = pa.Table.from_pylist(data_list)
        temp_file_path = f'/tmp/{uuid.uuid4()}.parquet'
        pq.write_table(table, temp_file_path)
        
        # 构造S3存储路径
        date_str = datetime.now().strftime('%Y-%m-%d')
        s3_key = f'user-data/{user_id}/{date_str}/{uuid.uuid4()}.parquet'
        
        # 上传至S3
        s3.upload_file(temp_file_path, BUCKET_NAME, s3_key)
    
    return {'statusCode': 200}

CDK v2部署配置

用CDK定义基础设施,重点配置DynamoDB Stream触发Lambda的批处理参数及权限:

CDK代码片段(Python示例)

from aws_cdk import (
    Stack,
    aws_dynamodb as dynamodb,
    aws_lambda as _lambda,
    aws_iam as iam,
    aws_lambda_event_sources as lambda_events,
    Duration,
)
from constructs import Construct

class DynamoDbToParquetStack(Stack):
    def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None:
        super().__init__(scope, construct_id, **kwargs)
        
        # 1. 创建DynamoDB表并启用Stream
        dynamo_table = dynamodb.Table(
            self, 'UserTable',
            partition_key=dynamodb.Attribute(name='id', type=dynamodb.AttributeType.STRING),
            stream=dynamodb.StreamViewType.NEW_IMAGE,  # 仅捕获新条目数据
            billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST
        )
        
        # 2. 创建Lambda函数
        parquet_converter_lambda = _lambda.Function(
            self, 'ParquetConverterLambda',
            runtime=_lambda.Runtime.PYTHON_3_11,
            handler='lambda_function.lambda_handler',
            code=_lambda.Code.from_asset('lambda'),  # Lambda代码所在目录
            memory_size=1024,  # 足够处理Parquet生成的内存
            timeout=Duration.seconds(90),  # 适配批量处理时长
            environment={
                'BUCKET_NAME': 'your-target-bucket'
            }
        )
        
        # 3. 配置DynamoDB Stream触发Lambda
        parquet_converter_lambda.add_event_source(
            lambda_events.DynamoEventSource(
                dynamo_table,
                starting_position=_lambda.StartingPosition.LATEST,
                batch_size=500,  # 每次触发最多处理500条
                batch_window=Duration.seconds(60),  # 最多等待60秒攒批
                retry_attempts=3
            )
        )
        
        # 4. 配置Lambda权限
        parquet_converter_lambda.add_to_role_policy(
            iam.PolicyStatement(
                actions=['s3:PutObject'],
                resources=['arn:aws:s3:::your-target-bucket/*']
            )
        )
        
        # 允许Lambda读取DynamoDB Stream
        dynamo_table.grant_stream_read(parquet_converter_lambda)

缓冲问题进阶处理

若担心Lambda重启导致内存中缓冲数据丢失,或需更灵活的缓冲策略:

  • 临时DynamoDB缓冲表:创建以user_id为分区键的表,存储待写入的data列表和累计大小,Lambda每次处理流事件时先写入缓冲表,再通过EventBridge定时触发另一个Lambda(如每分钟一次)检查缓冲表,达到阈值则生成Parquet写入S3并清空缓冲
  • S3分段上传:若单用户数据量极大,可直接用S3 Multipart Upload分块写入,避免内存不足

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 21:55:17