【发布时间】:2021-09-08 14:18:05
【问题描述】:
我想加入一个大流和一个小得多的流。我想广播较小的流,然后将其连接到较大的流。
但是我不确定如何处理存储广播的模式以及在processElement method 中如何查找匹配的模式然后组合这两个元素。
编辑: 我已经设法使用以下 sn-p 制作了广播加入的原型。我改编了官方培训存储库中的正常连接:https://github.com/apache/flink-training/blob/release-1.13/rides-and-fares/src/solution/scala/org/apache/flink/training/solutions/ridesandfares/scala/RidesAndFaresSolution.scala
这似乎有效,但我不确定我的逻辑是否正确。
//The main function has been abbreviated for ease of reading
def main(){
val rides = env
.addSource(rideSourceOrTest(new TaxiRideGenerator()))
.filter { ride => ride.isStart }
// .keyBy { ride => ride.rideId }
val fares = env
.addSource(fareSourceOrTest(new TaxiFareGenerator()))
val broadcastStateDescriptor = new MapStateDescriptor[Long,TaxiFare]("fares_broadcast",classOf[Long],classOf[TaxiFare])
val faresBroadcast: BroadcastStream[TaxiFare] = fares
.broadcast(broadcastStateDescriptor)
val result: DataStream[(TaxiRide,TaxiFare)] = rides
.connect(faresBroadcast)
.process(new BroadcastJoin())
}
class BroadcastJoin() extends BroadcastProcessFunction[TaxiRide,TaxiFare,(TaxiRide,TaxiFare)]{//IN1, IN2, OUT。 That is, non broadcast stream type, broadcast stream type and output stream type
//Broadcast state descriptor
private lazy val broadcastStateDescriptor = new MapStateDescriptor[Long,TaxiFare]("fares_broadcast",classOf[Long],classOf[TaxiFare])
//Process the broadcast stream element, value is the broadcast stream element passed in, and the modifiable broadcast state can be obtained through CTX
override def processBroadcastElement(value: TaxiFare, ctx: BroadcastProcessFunction[TaxiRide,TaxiFare,(TaxiRide,TaxiFare)]#Context, out: Collector[(TaxiRide,TaxiFare)]): Unit = {
val broadcast_status: BroadcastState[Long,TaxiFare] = ctx.getBroadcastState(broadcastStateDescriptor)
if(broadcast_status.contains(value.rideId)){
broadcast_status.remove(value.rideId)
}
broadcast_status.put ( value.rideId , value) // add the broadcast stream element to the broadcast state, which will be saved in local memory
}
//Handle non broadcast stream elements. Value is the non broadcast stream element passed in. Only read-only broadcast status can be obtained through CTX
override def processElement(value: TaxiRide, ctx: BroadcastProcessFunction[TaxiRide,TaxiFare,(TaxiRide,TaxiFare)]#ReadOnlyContext, out: Collector[(TaxiRide,TaxiFare)]): Unit = {
//Read broadcast status
val broadcast_status: ReadOnlyBroadcastState[Long, TaxiFare] = ctx.getBroadcastState(broadcastStateDescriptor)
if(broadcast_status.contains(value.rideId)) {
val foundMatch = broadcast_status.get(value.rideId)
out.collect((value, foundMatch)) //Send out the desired results
}
}
}
【问题讨论】:
-
欢迎来到 StackOverflow。你试过什么了?你在哪里卡住了?
-
嗨,彼得!我刚刚用更有意义的信息更新了我的帖子。我可以使用广播方法成功加入两个流,但我想知道我的实现是否有任何问题
标签: scala apache-flink flink-streaming