【发布时间】:2021-09-07 01:43:02
【问题描述】:
- 我有一个 RichParallelSourceFunction(parallelism=1),它每 5 秒查询一次 MySQL,它发出一个时间戳列表,指示何时开始/停止写入接收器。
- 此时间戳广播到原始流(并行度=10)。
- 我已将 RichParallelSourceFunction 的并行度设置为 1,以减少对 MySQL 服务器的同时请求数。
- 我很困惑在这种情况下是否需要广播状态。为什么不将广播数据存储在运营商本地数据结构中?
- .broadcast(stateDesc) 与 .broadcast() 之间有什么区别?
class MyBroadcastProcessFunction(name: String) extends BroadcastProcessFunction[Log, TimestampList, Log] with CheckpointedFunction {
private var sortedTimestamps: IndexedSeq[Long] = _
@transient var buffer: ListState[(Log, Long)] = _
def processElement(value: Log, timestamp: Long, out: Collector[Log]): Unit = {
if (timestamp > sortedTimestamps(0))
out.collect(value)
else
buffer.put(value, 0L)
}
override def processBroadcastElement(timestamps: TimestampList,
ctx: BroadcastProcessFunction[Log, TimestampList, Log]#Context,
out: Collector[Log]): Unit = {
// add timestamps to SortedTimestamps
// do I need to use BroadcastState mapstate here? or just use operator-local data structures (ex. sortedTimestamps)?
}
override def snapshotState(context: FunctionSnapshotContext): Unit = {}
override def initializeState(context: FunctionInitializationContext): Unit = {
val stateDesc = new ListStateDescriptor[(Log, Long)]("logBuffer",
classOf[(Log, Long)])
buffer = context.getOperatorStateStore.getListState(stateDesc)
}
}
【问题讨论】:
标签: scala apache-flink