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

如何将DynamoDB作为AWS Data Analytics(Flink)应用的Sink使用?

由于AWS Flink Data Analytics没有官方提供现成的DynamoDB Sink实现,你可以通过自定义Flink Sink函数来实现,主要有两种方式:同步Sink(适合小流量场景)和异步Sink(生产环境推荐,性能更高)。

一、自定义同步DynamoDB Sink

基于Flink的RichSinkFunction实现,每条数据同步写入DynamoDB。

1. 依赖准备

在你的项目中引入AWS SDK for DynamoDB的依赖(以Maven为例):

<dependency>
    <groupId>software.amazon.awssdk</groupId>
    <artifactId>dynamodb</artifactId>
    <version>2.20.0</version>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-streaming-java</artifactId>
    <version>1.17.0</version> <!-- 匹配你的Flink版本 -->
    <scope>provided</scope>
</dependency>

2. 实现Sink类

import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.dynamodb.DynamoDbClient;
import software.amazon.awssdk.services.dynamodb.model.PutItemRequest;
import software.amazon.awssdk.services.dynamodb.model.AttributeValue;
import java.util.Map;

public class DynamoDBSink<T> extends RichSinkFunction<T> {
    private transient DynamoDbClient dynamoDbClient;
    private final String tableName;
    private final Region region;

    public DynamoDBSink(String tableName, String regionId) {
        this.tableName = tableName;
        this.region = Region.of(regionId);
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        // 初始化DynamoDB客户端,AWS环境下自动通过IAM角色获取凭证
        dynamoDbClient = DynamoDbClient.builder()
                .region(region)
                .credentialsProvider(DefaultCredentialsProvider.create())
                .build();
    }

    @Override
    public void invoke(T value, Context context) throws Exception {
        // 将业务数据转换为DynamoDB的属性映射,需根据你的数据结构自定义实现
        Map<String, AttributeValue> item = convertToItemAttributes(value);
        
        PutItemRequest request = PutItemRequest.builder()
                .tableName(tableName)
                .item(item)
                .build();
        dynamoDbClient.putItem(request);
    }

    @Override
    public void close() throws Exception {
        if (dynamoDbClient != null) {
            dynamoDbClient.close();
        }
        super.close();
    }

    // 自定义数据转换逻辑,示例:将User对象转为DynamoDB Item
    private Map<String, AttributeValue> convertToItemAttributes(T value) {
        // User user = (User) value;
        // return Map.of(
        //     "userId", AttributeValue.builder().s(user.getUserId()).build(),
        //     "userName", AttributeValue.builder().s(user.getUserName()).build(),
        //     "age", AttributeValue.builder().n(String.valueOf(user.getAge())).build()
        // );
        throw new UnsupportedOperationException("请根据你的数据结构实现转换逻辑");
    }
}

3. 在Flink作业中使用

// 假设stream是你的数据流
DataStream<User> userStream = ...;
userStream.addSink(new DynamoDBSink<>("user-table", "us-east-1"));

二、自定义异步DynamoDB Sink(生产环境推荐)

同步Sink会阻塞数据流,导致性能瓶颈。使用Flink的AsyncFunction结合DynamoDB异步客户端实现非阻塞写入,大幅提升吞吐量。

1. 实现异步Sink类

import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.async.ResultFuture;
import org.apache.flink.streaming.api.functions.async.RichAsyncFunction;
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClient;
import software.amazon.awssdk.services.dynamodb.model.PutItemRequest;
import software.amazon.awssdk.services.dynamodb.model.AttributeValue;
import java.util.Collections;
import java.util.Map;
import java.util.concurrent.CompletableFuture;

public class AsyncDynamoDBSink<T> extends RichAsyncFunction<T, Void> {
    private transient DynamoDbAsyncClient asyncClient;
    private final String tableName;
    private final Region region;

    public AsyncDynamoDBSink(String tableName, String regionId) {
        this.tableName = tableName;
        this.region = Region.of(regionId);
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        asyncClient = DynamoDbAsyncClient.builder()
                .region(region)
                .credentialsProvider(DefaultCredentialsProvider.create())
                .build();
    }

    @Override
    public void asyncInvoke(T value, ResultFuture<Void> resultFuture) throws Exception {
        Map<String, AttributeValue> item = convertToItemAttributes(value);
        PutItemRequest request = PutItemRequest.builder()
                .tableName(tableName)
                .item(item)
                .build();

        // 异步提交写入请求,处理回调
        asyncClient.putItem(request)
                .thenAccept(response -> resultFuture.complete(Collections.emptyList()))
                .exceptionally(ex -> {
                    resultFuture.completeExceptionally(ex);
                    return null;
                });
    }

    @Override
    public void close() throws Exception {
        if (asyncClient != null) {
            asyncClient.close();
        }
        super.close();
    }

    // 自定义数据转换逻辑,同同步Sink
    private Map<String, AttributeValue> convertToItemAttributes(T value) {
        throw new UnsupportedOperationException("请根据你的数据结构实现转换逻辑");
    }
}

2. 在Flink作业中使用

import org.apache.flink.streaming.api.async.AsyncDataStream;
import java.util.concurrent.TimeUnit;

DataStream<User> userStream = ...;
AsyncDataStream.unorderedWait(
        userStream,
        new AsyncDynamoDBSink<>("user-table", "us-east-1"),
        1000, // 单个请求超时时间(毫秒)
        TimeUnit.MILLISECONDS,
        100 // 并发请求数,根据DynamoDB吞吐量调整
);

三、最佳实践

  • 批量写入优化:攒一批数据后使用BatchWriteItemRequest批量提交,减少API调用次数。可在Sink中维护缓冲区,达到阈值或定时触发批量写入。
  • 权限配置:给Flink应用的执行IAM角色添加DynamoDB写入权限,推荐最小权限策略(如仅允许dynamodb:PutItem、dynamodb:BatchWriteItem)。
  • 错误处理:针对DynamoDB限流错误(ProvisionedThroughputExceededException)添加指数退避重试;开启Flink Checkpoint保证数据不丢失。
  • 客户端调优:调整DynamoDB异步客户端的连接池大小、超时时间,适配Flink作业的并行度和数据量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 12:31:03