Apache Flink: kafkasource с холодным холодильником работает плохо при подключении вещанияJAVA

Программисты JAVA общаются здесь
Anonymous
Apache Flink: kafkasource с холодным холодильником работает плохо при подключении вещания

Сообщение Anonymous »

Когда Kafkasource Connected A VroadcastStream устанавливается с бездельностью, водяной знак ниже по течению является ненормальным. < /p>
Мой вопрос: как сделать водяной маркирной для использования в окне. Подписывается < /p>
zkwatchersource установлен с WatermarkStrategy, который всегда отправляет watermarm.max_watermark, чтобы продвинуть нисходящий водяной знак. public static class MaxWatermarkGenerator implements WatermarkGenerator {

@Override
public void onEvent(T event, long eventTimestamp, WatermarkOutput output) {
}

@Override
public void onPeriodicEmit(WatermarkOutput output) {
output.emitWatermark(Watermark.MAX_WATERMARK);
}
}

Следующим, водяной знак оператора1 является long.min_value сначала, но через некоторое время он становится до Long.max_value
Watermark window operator2 water.max_value навсегда, окно не может работать нормально
watermark
public class KafkaSourceWithBroadcastConnectDemo {
public static void main(String[] args) throws Exception {

StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());
env.setParallelism(2);
Configuration conf = FlinkUtils.loadConfig(args);

// kafka source
DataStreamSource kafkaSourceStream = env.fromSource(
Sources.kafkaSourceBuilder()
.setBootstrapServers(conf.get(JobOptions.MAIN_KAFKA_BOOTSTRAP_SERVERS))
.setTopicPattern(Pattern.compile(conf.get(JobOptions.KAFKA_PERFORMANCE_TOPIC_PATTERN)))
.setGroupId(conf.get(JobOptions.KAFKA_CONSUMER_GROUP_ID))
.setValueOnlyDeserializer(new JacksonDeserializationSchema(PerformanceMessage.class))
.build(),
WatermarkStrategy
.forBoundedOutOfOrderness(Duration.ofSeconds(conf.get(JobOptions.WATERMARK_MAX_OUT_OF_ORDERNESS)))
.withIdleness(Duration.ofSeconds(conf.get(JobOptions.KAFKA_SOURCE_IDLE_TIMEOUT)))
.withTimestampAssigner((msg, timestamp) -> NumberUtils.parseLong(msg.getTime())),
"kafka-source")
.setParallelism(1);

// broadcast source
BroadcastStream broadcastStream = env.fromSource(
new DataGeneratorSource((GeneratorFunction) value -> "number:" + value, Long.MAX_VALUE, RateLimiterStrategy.perSecond(1), Types.STRING),
WatermarkStrategy
.forGenerator(ctx -> new MaxWatermarkGenerator()),
"max-set").broadcast(CONTROL_SIGNAL_BROADCAST);

// connect and process
kafkaSourceStream
.connect(broadcastStream)
.process(new BroadcastProcessFunction() {
@Override
public void processElement(PerformanceMessage value, BroadcastProcessFunction.ReadOnlyContext ctx, Collector out) throws Exception {
System.err.println(ctx.currentWatermark() + "--->" + value);
out.collect(value);
}
@Override
public void processBroadcastElement(String value, BroadcastProcessFunction.Context ctx, Collector out) throws Exception {
// do nothing
}
}).setParallelism(2)
.keyBy(msg -> msg.getDevice().getId())
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(1L)))
.process(new ProcessWindowFunction() {
@Override
public void process(String deviceNo, ProcessWindowFunction.Context context, Iterable elements, Collector out) throws Exception {
System.out.println(context.currentWatermark());
for (PerformanceMessage message : elements) {
out.collect(message.getDevice().getId() + ":" + message.getTime());
}
}
}).print().setParallelism(4);

env.execute();
}

public static final MapStateDescriptor CONTROL_SIGNAL_BROADCAST =
new MapStateDescriptor(
"control-signal-broadcast",
BasicTypeInfo.STRING_TYPE_INFO,
BasicTypeInfo.VOID_TYPE_INFO
);

public static class MaxWatermarkGenerator implements WatermarkGenerator {

@Override
public void onEvent(T event, long eventTimestamp, WatermarkOutput output) {
}

@Override
public void onPeriodicEmit(WatermarkOutput output) {
output.emitWatermark(Watermark.MAX_WATERMARK);
}
}
}


Подробнее здесь: https://stackoverflow.com/questions/797 ... caststream

Вернуться в «JAVA»