求助:Apache DataStream集成StateFun的远程有状态函数连接示例
Apache Flink StateFun与DataStream集成的可运行代码获取方案
问题说明
- 学习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
相关产品推荐
相关产品推荐

