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

求助:Apache DataStream集成StateFun的远程有状态函数连接示例

问题说明

  • 学习StateFun与DataStream集成时,参考文档提供的完整示例代码链接出现404错误;
  • 已找到的示例代码使用withFunctionProvider绑定嵌入式函数,与文档中的非嵌入式函数示例不符。

解决途径

1. 匹配文档版本的示例代码

你参考的文档基于特定commit(7ec6664),可通过以下方式获取对应版本的示例:

  • 克隆StateFun仓库后,执行git checkout 7ec6664切换到该版本分支;
  • 进入statefun-examples/statefun-flink-datastream-example目录,即可获取与文档完全匹配的非嵌入式函数集成代码。

2. 使用最新稳定版官方示例

若无需严格匹配旧版本,直接克隆StateFun官方仓库,查看statefun-examples模块下的statefun-flink-datastream-example:

  • 该目录下的代码对应最新稳定版,包含远程函数(非嵌入式)与DataStream集成的完整实现,可直接运行测试。

3. 基于文档片段手动补全示例

参考文档中的代码片段,结合StateFun API规范,可快速补全非嵌入式函数的集成代码,示例如下:

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.statefun.flink.datastream.StatefulFunctionDataStreamBuilder;
import org.apache.flink.statefun.flink.datastream.StatefulFunctionEgressStreams;
import org.apache.flink.statefun.sdk.FunctionType;
import org.apache.flink.statefun.sdk.io.EgressIdentifier;
import org.apache.flink.statefun.sdk.io.IngressIdentifier;
import org.apache.flink.statefun.sdk.reqreply.generated.RequestReply;

import java.net.URI;
import java.time.Duration;

public class StateFunDataStreamRemoteExample {
    // 定义函数类型、Ingress/Egress标识
    private static final FunctionType REMOTE_GREET_FUNCTION = new FunctionType("example", "greet");
    private static final IngressIdentifier<String> NAMES_INGRESS = new IngressIdentifier<>(String.class, "example", "names");
    private static final EgressIdentifier<RequestReply.IngressMessage> GREETINGS_EGRESS = 
        new EgressIdentifier<>(RequestReply.IngressMessage.class, "example", "greetings");

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 模拟输入数据流
        DataStream<String> namesInput = env.fromElements("Alice", "Bob", "Charlie");

        // 构建StateFun与DataStream的集成,仅使用远程函数
        StatefulFunctionEgressStreams egressStreams = StatefulFunctionDataStreamBuilder.builder("example")
                .withDataStreamAsIngress(NAMES_INGRESS, namesInput)
                .withRequestReplyRemoteFunction(
                        StatefulFunctionDataStreamBuilder.requestReplyFunctionBuilder(
                                REMOTE_GREET_FUNCTION, URI.create("http://localhost:5000/statefun"))
                                .withMaxRequestDuration(Duration.ofSeconds(15))
                                .withMaxNumBatchRequests(500))
                .withEgressId(GREETINGS_EGRESS)
                .build(env);

        // 输出Egress结果
        egressStreams.getEgress(GREETINGS_EGRESS).print();

        env.execute("StateFun DataStream Remote Integration");
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 05:15:35