两个Flink窗口流输出至同一Kinesis Sink无数据,如何解决?
问题
将两个Flink窗口流输出至同一个Kinesis Sink时,无任何结果写入该Sink;移除其中一个窗口流后,结果可正常发布;同时添加两个流至该Sink时,二者输出均失效。如何实现让两个窗口流的结果都输出到同一个Kinesis Sink?
相关代码:
public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); ObjectMapper jsonParser = new ObjectMapper(); DataStream<String> inputStream = createKinesisSource(env); FlinkKinesisProducer<String> kinesisSink = createKinesisSink(); WindowedStream oneMinStream = inputStream .map(value -> jsonParser.readValue(value, JsonNode.class)) .keyBy(node -> node.get("accountId")) .window(TumblingProcessingTimeWindows.of(Time.minutes(1))); oneMinStream .aggregate(new LoginAggregator("k1m")) .addSink(kinesisSink); WindowedStream twoMinStream = inputStream .map(value -> jsonParser.readValue(value, JsonNode.class)) .keyBy(node -> node.get("accountId")) .window(TumblingProcessingTimeWindows.of(Time.minutes(2))); twoMinStream .aggregate(new LoginAggregator("k2m")) .addSink(kinesisSink); try { env.execute("Flink Kinesis Streaming Sink Job"); } catch (Exception e) { LOG.error("failed"); LOG.error(e.getLocalizedMessage()); LOG.error(e.getStackTrace().toString()); throw e; } } private static DataStream<String> createKinesisSource(StreamExecutionEnvironment env) { Properties inputProperties = new Properties(); inputProperties.setProperty(ConsumerConfigConstants.AWS_REGION, region); inputProperties.setProperty(ConsumerConfigConstants.STREAM_INITIAL_POSITION, "LATEST"); return env.addSource(new FlinkKinesisConsumer<>(inputStreamName, new SimpleStringSchema(), inputProperties)); } private static FlinkKinesisProducer<String> createKinesisSink() { Properties outputProperties = new Properties(); outputProperties.setProperty(ConsumerConfigConstants.AWS_REGION, region); outputProperties.setProperty("AggregationEnabled", "false"); FlinkKinesisProducer<String> sink = new FlinkKinesisProducer<>(new SimpleStringSchema(), outputProperties); sink.setDefaultStream(outputStreamName); sink.setDefaultPartition(UUID.randomUUID().toString()); return sink; }
解决方案
核心原因:FlinkKinesisProducer是线程不安全的,不能被多个数据流共享使用。复用同一个sink实例会导致内部状态混乱,最终无法正常写入Kinesis。
以下两种方法可解决问题:
方法1:为每个流创建独立的Kinesis Sink实例
直接修改代码,让每个流的addSink都使用新创建的FlinkKinesisProducer实例,避免复用同一个对象:
public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); ObjectMapper jsonParser = new ObjectMapper(); DataStream<String> inputStream = createKinesisSource(env); WindowedStream oneMinStream = inputStream .map(value -> jsonParser.readValue(value, JsonNode.class)) .keyBy(node -> node.get("accountId")) .window(TumblingProcessingTimeWindows.of(Time.minutes(1))); // 每个流单独创建Sink实例 oneMinStream .aggregate(new LoginAggregator("k1m")) .addSink(createKinesisSink()); WindowedStream twoMinStream = inputStream .map(value -> jsonParser.readValue(value, JsonNode.class)) .keyBy(node -> node.get("accountId")) .window(TumblingProcessingTimeWindows.of(Time.minutes(2))); twoMinStream .aggregate(new LoginAggregator("k2m")) .addSink(createKinesisSink()); try { env.execute("Flink Kinesis Streaming Sink Job"); } catch (Exception e) { LOG.error("failed"); LOG.error(e.getLocalizedMessage()); LOG.error(e.getStackTrace().toString()); throw e; } }
createKinesisSink方法无需修改,每次调用都会返回全新的Sink实例。
方法2:合并两个窗口流后统一输出(可选)
如果两个窗口流的输出格式一致,可先将两个流合并为一个数据流,再统一输出到单个Sink,减少资源占用:
public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); ObjectMapper jsonParser = new ObjectMapper(); DataStream<String> inputStream = createKinesisSource(env); WindowedStream oneMinStream = inputStream .map(value -> jsonParser.readValue(value, JsonNode.class)) .keyBy(node -> node.get("accountId")) .window(TumblingProcessingTimeWindows.of(Time.minutes(1))); DataStream<String> oneMinResult = oneMinStream.aggregate(new LoginAggregator("k1m")); WindowedStream twoMinStream = inputStream .map(value -> jsonParser.readValue(value, JsonNode.class)) .keyBy(node -> node.get("accountId")) .window(TumblingProcessingTimeWindows.of(Time.minutes(2))); DataStream<String> twoMinResult = twoMinStream.aggregate(new LoginAggregator("k2m")); // 合并两个结果流 DataStream<String> mergedStream = oneMinResult.union(twoMinResult); // 输出到单个Sink mergedStream.addSink(createKinesisSink()); try { env.execute("Flink Kinesis Streaming Sink Job"); } catch (Exception e) { LOG.error("failed"); LOG.error(e.getLocalizedMessage()); LOG.error(e.getStackTrace().toString()); throw e; } }
内容的提问来源于stack exchange,提问作者Thor
相关产品推荐
相关产品推荐

