Akka Stream中divertTo/alsoTo算子的物化值类型保留及自定义Flow物化值类型实现咨询
如何为Akka Streams Flow/Sink指定自定义物化值类型
当然可以创建带有预定义物化值类型的Flow(或Sink)——这正是Akka Streams中*Mat系列方法的设计目的!你之前的思路方向是对的,但需要调整一下实现方式,尤其是结合Kafka消费者的Consumer.DrainingControl场景,这里有几个关键点要注意:
核心问题解析
你当前用Flow.to(Committer.sink(...))得到的Sink物化值是NotUsed,因为不带Mat的方法(比如to、via)会默认丢弃上游/下游的物化值,强制返回NotUsed。要保留自定义物化值,必须使用toMat、viaMat这类带物化值控制的方法,并通过Keep策略指定要保留的物化值。
另外,Consumer.DrainingControl需要结合**Source的Consumer.Control和Committer的Committer.Control**来创建——你没办法在Sink内部单独生成它,因为Sink无法访问Source的Consumer.Control,必须在流的顶层把这两个控制对象组合起来。
具体实现方案
1. 构建带自定义物化值的Flow
先创建一个Flow,它处理消息的同时,保留Committer的控制对象作为物化值:
import akka.stream.javadsl.Flow; import akka.stream.javadsl.Keep; import akka.kafka.Committer; import akka.kafka.ConsumerMessage; import scala.Tuple2; // 这个Flow接收Kafka的CommittableMessage,转换为SomeType,物化值是Committer.Control private Flow<ConsumerMessage.CommittableMessage<String, String>, SomeType, Committer.Control> createProcessingFlow(CommitterSettings committerSettings) { return Flow.of(ConsumerMessage.CommittableMessage.class) // 转换消息并保留可提交的Offset .map(msg -> new Tuple2<>(convertToSomeType(msg.record()), msg.committableOffset())) // 分流到Committer Sink,保留Committer的Control作为Flow的物化值 .alsoToMat(Committer.sink(committerSettings), Keep.right()) // 只输出处理后的业务对象 .map(Tuple2::_1); } // 辅助方法:把Kafka记录转换为你的业务类型 private SomeType convertToSomeType(ConsumerMessage.Record<String, String> record) { // 实现你的消息转换逻辑 return new SomeType(record.value()); }
2. 组合流并生成DrainingControl
接下来,在流的顶层,结合Source的Consumer.Control和Flow的Committer.Control来创建Consumer.DrainingControl,这里用divertToMat来处理分流场景:
import akka.kafka.Consumer; import akka.kafka.ConsumerSettings; import akka.kafka.Subscriptions; import akka.stream.javadsl.Sink; import akka.stream.javadsl.RunnableGraph; // 初始化Kafka消费者和Committer配置 ConsumerSettings<String, String> consumerSettings = ConsumerSettings.create(system, stringDeserializer, stringDeserializer) .withBootstrapServers("localhost:9092") .withGroupId("my-consumer-group"); CommitterSettings committerSettings = CommitterSettings.create(system); // 创建Kafka Source,物化值是Consumer.Control Source<ConsumerMessage.CommittableMessage<String, String>, Consumer.Control> kafkaSource = Consumer.plainSource(consumerSettings, Subscriptions.topics("my-topic")); // 主处理路径的Sink(示例用忽略Sink,替换成你的实际处理逻辑) Sink<SomeType, NotUsed> mainSink = Sink.ignore(); // 构建分流的Flow+Sink组合 Sink<ConsumerMessage.CommittableMessage<String, String>, Committer.Control> divertSink = createProcessingFlow(committerSettings).toMat(mainSink, Keep.left()); // 使用divertToMat组合物化值,生成DrainingControl RunnableGraph<Consumer.DrainingControl> runnableGraph = kafkaSource .divertToMat( divertSink, // 把Source的Consumer.Control和分流Sink的Committer.Control组合成DrainingControl (consumerControl, committerControl) -> Consumer.createDrainingControl(consumerControl, committerControl) ) // 主路径可以是任意处理,这里用空Sink示例 .to(Sink.ignore()); // 运行流,得到DrainingControl用于优雅关闭 Consumer.DrainingControl drainingControl = runnableGraph.run(system); // 后续可以调用drainingControl.drainAndShutdown()来优雅停止消费者和提交器
关键要点
- 使用
*Mat方法:所有需要保留物化值的操作,都用带Mat后缀的方法(toMat、viaMat、divertToMat),避免默认丢弃物化值的普通方法。 Keep策略:用Keep.right()保留下游物化值,Keep.left()保留上游,Keep.both()保留两者(会返回Tuple2),根据你的需求选择。- DrainingControl的生成时机:必须在流的顶层组合
Consumer.Control和Committer.Control,因为Sink/Flow无法访问Source的控制对象。
内容的提问来源于stack exchange,提问作者Yannic Bürgmann
相关产品推荐
相关产品推荐

