You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

两个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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 12:15:33