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

如何为AWS SageMaker Feature Store实现Apache Flink自定义Sink?

实现Apache Flink自定义Sink对接AWS SageMaker Feature Store

基类选择建议

先明确两类基类的适用场景,帮你快速做决策:

  • RichSinkFunction/SinkFunction:传统Sink实现方式,上手简单,适合快速验证或低吞吐量场景,但需要自行实现批量、异步、重试、容错等逻辑。
  • AsyncSinkBase/AsyncSinkWriter:Flink新一代连接器框架(Connector Base)的组件,内置异步处理、批量提交、至少一次语义等生产级能力,和SourceReaderBase属于同一套标准化框架,优先推荐用于生产环境。

选择AsyncSinkBase的核心原因:它已经封装了至少一次语义的容错逻辑,且支持异步化同步调用(提升吞吐量),不用从零搭建重试、 checkpoint 对齐等机制。

基于AsyncSinkBase的生产级实现步骤

1. 依赖准备

引入必要的Flink和AWS SDK依赖(以Maven为例):

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-base</artifactId>
    <version>${flink.version}</version>
</dependency>
<dependency>
    <groupId>software.amazon.awssdk</groupId>
    <artifactId>sagemakerfeaturestoreruntime</artifactId>
    <version>${aws.sdk.version}</version>
</dependency>

2. 核心组件实现

(1)实现异步SinkWriter

继承AsyncSinkWriter,封装SageMaker Feature Store的写入逻辑:

public class SageMakerFeatureStoreSinkWriter extends AsyncSinkWriter<YourRecordType, Void> {
    private final AmazonSageMakerFeatureStoreRuntime client;
    private final String featureGroupName;

    public SageMakerFeatureStoreSinkWriter(
            AmazonSageMakerFeatureStoreRuntime client,
            String featureGroupName,
            AsyncSinkWriterBuilder<YourRecordType, Void> builder) {
        super(builder);
        this.client = client;
        this.featureGroupName = featureGroupName;
    }

    @Override
    protected CompletableFuture<Void> sendRequest(YourRecordType record) {
        // 将Flink数据转换为SageMaker PutRecord请求
        PutRecordRequest request = PutRecordRequest.builder()
                .featureGroupName(featureGroupName)
                .record(convertToFeatureValues(record))
                .build();
        // 用CompletableFuture包装同步调用,实现异步化
        return CompletableFuture.runAsync(() -> client.putRecord(request));
    }

    // 自定义数据转换逻辑:将业务对象转为FeatureValue列表
    private List<FeatureValue> convertToFeatureValues(YourRecordType record) {
        List<FeatureValue> featureValues = new ArrayList<>();
        featureValues.add(FeatureValue.builder()
                .featureName("record_id")
                .valueAsString(record.getRecordId())
                .build());
        // 添加其他字段转换逻辑
        return featureValues;
    }
}

(2)实现Sink主体

继承AsyncSinkBase,构建可配置的Sink实例:

public class SageMakerFeatureStoreSink extends AsyncSinkBase<YourRecordType, Void> {
    private final String featureGroupName;
    private final AwsCredentialsProvider credentialsProvider;
    private final Region region;

    public SageMakerFeatureStoreSink(
            String featureGroupName,
            AwsCredentialsProvider credentialsProvider,
            Region region) {
        this.featureGroupName = featureGroupName;
        this.credentialsProvider = credentialsProvider;
        this.region = region;
    }

    @Override
    public AsyncSinkWriter<YourRecordType, Void> createWriter(AsyncSinkWriterBuilder<YourRecordType, Void> builder) {
        // 初始化SageMaker客户端(避免序列化问题,在创建Writer时初始化)
        AmazonSageMakerFeatureStoreRuntime client = AmazonSageMakerFeatureStoreRuntimeClient.builder()
                .credentialsProvider(credentialsProvider)
                .region(region)
                .build();
        return new SageMakerFeatureStoreSinkWriter(client, featureGroupName, builder);
    }

    // 提供Builder模式,方便Flink作业中配置使用
    public static class Builder {
        private String featureGroupName;
        private AwsCredentialsProvider credentialsProvider = DefaultCredentialsProvider.create();
        private Region region;

        public Builder featureGroupName(String featureGroupName) {
            this.featureGroupName = featureGroupName;
            return this;
        }

        public Builder credentialsProvider(AwsCredentialsProvider credentialsProvider) {
            this.credentialsProvider = credentialsProvider;
            return this;
        }

        public Builder region(Region region) {
            this.region = region;
            return this;
        }

        public SageMakerFeatureStoreSink build() {
            return new SageMakerFeatureStoreSink(featureGroupName, credentialsProvider, region);
        }
    }
}

3. 确保至少一次语义

AsyncSinkBase默认支持至少一次语义:

  • Checkpoint触发时,会等待所有异步写入请求完成后再完成checkpoint
  • 可通过配置重试参数(如setRetryStrategy)处理写入失败的情况
  • 结合SageMaker Feature Store底层DynamoDB的键值特性,只要RecordIdentifier唯一,重复写入不会产生副作用,也可基于此实现恰好一次语义

基于RichSinkFunction的轻量实现方案

如果是快速验证场景,可使用RichSinkFunction快速实现:

public class SageMakerFeatureStoreRichSink extends RichSinkFunction<YourRecordType> {
    private transient AmazonSageMakerFeatureStoreRuntime client;
    private final String featureGroupName;
    private final String region;

    public SageMakerFeatureStoreRichSink(String featureGroupName, String region) {
        this.featureGroupName = featureGroupName;
        this.region = region;
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        // 在open方法初始化客户端,避免序列化问题
        client = AmazonSageMakerFeatureStoreRuntimeClient.builder()
                .region(Region.of(region))
                .build();
    }

    @Override
    public void invoke(YourRecordType record, Context context) throws Exception {
        PutRecordRequest request = PutRecordRequest.builder()
                .featureGroupName(featureGroupName)
                .record(convertToFeatureValues(record))
                .build();
        // 添加重试逻辑,确保至少一次语义
        int retryTimes = 3;
        while (retryTimes-- > 0) {
            try {
                client.putRecord(request);
                break;
            } catch (SdkException e) {
                if (retryTimes == 0) throw e;
                Thread.sleep(1000);
            }
        }
    }

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

    private List<FeatureValue> convertToFeatureValues(YourRecordType record) {
        // 数据转换逻辑
    }
}

注意:该方案为同步写入,吞吐量较低,需自行实现批量、重试等优化逻辑。

关键注意事项

  • 客户端初始化:Flink算子会被序列化,不能将AWS客户端直接作为成员变量,需在open(RichSinkFunction)或SinkWriter创建时初始化。
  • 异步优化:即使使用同步客户端,也要用CompletableFuture包装调用,避免阻塞Flink任务线程。
  • 限流与重试:配置合理的重试次数和间隔,应对SageMaker的限流或网络波动。
  • 批量写入:若需更高吞吐量,可攒批后调用批量接口(若支持),或在AsyncSinkWriter中配置批量提交的大小和超时时间。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 03:46:08