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

Flink无键流数据如何使用滚动窗口(Tumbling Window)提升处理吞吐量

问题描述

我需要在处理无键流数据的程序中使用滚动窗口(Tumbling Window)功能,当前我的流处理吞吐量仅为300条/秒,期望至少提升到5000条/秒,因此计划设置2秒的滚动窗口来优化性能,但不清楚该场景下具体如何实现。
注意:我使用Geomesa HBase平台存储消息,此处仅贴出与窗口函数实现相关的代码片段,未粘贴全部应用代码,足够支撑窗口功能的需求理解
我的现有Flink代码如下:

public class Tranport {
    public static void main(String[] args) throws Exception {
        // fetch runtime arguments
        String bootstrapServers = "xx.xx.xxx.xxx:xxxx";
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Set up the Consumer and create a datastream from this source
        Properties properties = new Properties();
        properties.setProperty("bootstrap.servers", bootstrapServers);
        properties.setProperty("group.id", "group_id");
        final FlinkKafkaConsumer<String> flinkConsumer = new FlinkKafkaConsumer<>("lc", new SimpleStringSchema(), properties);
        flinkConsumer.setStartFromTimestamp(Long.parseLong("0"));
        DataStream<String> readingStream = env.addSource(flinkConsumer);
        readingStream.rebalance().map(new RichMapFunction<String, String>() {
            private static final long serialVersionUID = -2547861355L; // random number
            DataStore lc_live = null;
            SimpleFeatureType sft_live;
            SimpleFeatureBuilder SFbuilderLive; // feature builder for live
            List<SimpleFeature> lc_live_features; // 
            @Override
            public void open(Configuration parameters) throws Exception {
                System.out.println("In open method.");
                // --- GEOMESA, GEOTOOLS APPROACH ---//
                // define connection parameters to xxx GeoMesa-HBase DataStore
                Map<String, Serializable> params_live = new HashMap<>();
                params_live.put("xxxx", "xxx"); // HBase table name
                params_live.put("xxxx","xxxx");
                try {
                    lc_live = DataStoreFinder.getDataStore(params_live);
                    if (lc_live == null) {
                        System.out.println("Could not connect to live");
                    } else {
                        System.out.println("Successfully connected to live");
                    }
                } catch (IOException e) {
                    e.printStackTrace();
                }
                // create simple feature type for x table in HBASE 
                StringBuilder attributes1 = new StringBuilder();
                attributes1.append("xxx:String,");
                attributes1.append("xxx:Long,");
                attributes1.append("source:String,");
                attributes1.append("xxx:String,");
                attributes1.append("xxx:Double,");
                attributes1.append("status:String,");
                attributes1.append("forecast:Double,");
                attributes1.append("carsCount:Integer,");
                attributes1.append("*xxx:Point:srid=4326");
                sft_history = SimpleFeatureTypes.createType("xxxx", attributes1.toString());
                try {
                    lc_history.createSchema(sft_history);                                               
                } catch (IOException e) {
                    e.printStackTrace();
                }
                // Initialize the variables
                numberOfMessagesProcessed = 0;
                numberOfMessagesFailed = 0;
                numberOfMessagesSkipped = 0;
        // for lc_Live
                lc_live_features = new ArrayList<>();
                SFbuilderLive = new SimpleFeatureBuilder(sft_live);

我需要在此处创建一个滚动窗口函数(Window All),收集2秒窗口内的所有流消息,写入我下方定义的ArrayList中,用于后续批量写入存储

// live GeoMesa-HBase DataStore
                        // copy the list into a local variable and empty the list for the next iteration
                        List<SimpleFeature> LocalFeatures = live_features;
                        live_features = new ArrayList<>();
                        LocalFeatures = Collections.unmodifiableList(LocalFeatures);
                        try (FeatureWriter<SimpleFeatureType, SimpleFeature> writer = live.getFeatureWriterAppend(sft_live.getTypeName(), Transaction.AUTO_COMMIT)) {
                            System.out.println("Writing " + LocalFeatures.size() + " features to live");
                            for (SimpleFeature feature : LocalFeatures) {
                                SimpleFeature toWrite = writer.next();
                                toWrite.setAttributes(feature.getAttributes());
                                ((FeatureIdImpl) toWrite.getIdentifier()).setID(feature.getID());
                                toWrite.getUserData().put(Hints.USE_PROVIDED_FID, Boolean.TRUE);
                                toWrite.getUserData().putAll(feature.getUserData());
                                writer.write();
                            }
                        } catch (IOException e) {
                            e.printStackTrace();
                        }
实现方案

你当前的写法是在MapFunction里手动攒数据,缺少窗口触发逻辑,且无状态保障容易丢失数据,直接用Flink自带的无键滚动窗口机制即可满足需求,修改后的完整pipeline代码如下:

public class Tranport {
    public static void main(String[] args) throws Exception {
        String bootstrapServers = "xx.xx.xxx.xxx:xxxx";
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 可选:开启checkpoint保证数据不丢,按自己需求配置
        env.enableCheckpointing(30000);
        
        Properties properties = new Properties();
        properties.setProperty("bootstrap.servers", bootstrapServers);
        properties.setProperty("group.id", "group_id");
        final FlinkKafkaConsumer<String> flinkConsumer = new FlinkKafkaConsumer<>("lc", new SimpleStringSchema(), properties);
        flinkConsumer.setStartFromTimestamp(Long.parseLong("0"));
        DataStream<String> readingStream = env.addSource(flinkConsumer);
        
        // 第一步:将Kafka字符串解析为SimpleFeature
        DataStream<SimpleFeature> featureStream = readingStream.rebalance()
            .map(new RichMapFunction<String, SimpleFeature>() {
                private static final long serialVersionUID = -2547861355L;
                private transient SimpleFeatureBuilder sfBuilderLive;
                private transient SimpleFeatureType sftLive;
                @Override
                public void open(Configuration parameters) throws Exception {
                    // 初始化Schema和构造器
                    StringBuilder attributes1 = new StringBuilder();
                    attributes1.append("xxx:String,");
                    attributes1.append("xxx:Long,");
                    attributes1.append("source:String,");
                    attributes1.append("xxx:String,");
                    attributes1.append("xxx:Double,");
                    attributes1.append("status:String,");
                    attributes1.append("forecast:Double,");
                    attributes1.append("carsCount:Integer,");
                    attributes1.append("*xxx:Point:srid=4326");
                    sftLive = SimpleFeatureTypes.createType("xxxx", attributes1.toString());
                    sfBuilderLive = new SimpleFeatureBuilder(sftLive);
                }
                @Override
                public SimpleFeature map(String value) throws Exception {
                    sfBuilderLive.reset();
                    // 这里替换成你自己的字段解析逻辑,把kafka消息的字段逐一对sft的字段赋值
                    // 示例:sfBuilderLive.set("source", "kafka_lc");
                    // 示例:sfBuilderLive.set("carsCount", 10);
                    // 省略你自己的解析逻辑,最后返回构造好的SimpleFeature
                    return sfBuilderLive.buildFeature(null);
                }
            });
        
        // 第二步:开2秒处理时间滚动窗口,批量写入Geomesa
        featureStream
            .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(2)))
            .process(new ProcessAllWindowFunction<SimpleFeature, Void, TimeWindow>() {
                private transient DataStore lcLive;
                private transient SimpleFeatureType sftLive;
                @Override
                public void open(Configuration parameters) throws Exception {
                    // 初始化Geomesa连接,每个subtask初始化一次
                    Map<String, Serializable> paramsLive = new HashMap<>();
                    paramsLive.put("xxxx", "xxx"); // HBase表名
                    paramsLive.put("xxxx","xxxx");
                    lcLive = DataStoreFinder.getDataStore(paramsLive);
                    if (lcLive == null) {
                        throw new RuntimeException("Geomesa HBase连接失败");
                    }
                    // 初始化Schema
                    StringBuilder attributes1 = new StringBuilder();
                    attributes1.append("xxx:String,");
                    attributes1.append("xxx:Long,");
                    attributes1.append("source:String,");
                    attributes1.append("xxx:String,");
                    attributes1.append("xxx:Double,");
                    attributes1.append("status:String,");
                    attributes1.append("forecast:Double,");
                    attributes1.append("carsCount:Integer,");
                    attributes1.append("*xxx:Point:srid=4326");
                    sftLive = SimpleFeatureTypes.createType("xxxx", attributes1.toString());
                }
                @Override
                public void process(Context context, Iterable<SimpleFeature> elements, Collector<Void> out) throws Exception {
                    List<SimpleFeature> localFeatures = new ArrayList<>();
                    elements.forEach(localFeatures::add);
                    if (localFeatures.isEmpty()) {
                        return;
                    }
                    // 直接复用你原来的批量写入逻辑
                    try (FeatureWriter<SimpleFeatureType, SimpleFeature> writer = lcLive.getFeatureWriterAppend(sftLive.getTypeName(), Transaction.AUTO_COMMIT)) {
                        System.out.println("Writing " + localFeatures.size() + " features to live");
                        for (SimpleFeature feature : localFeatures) {
                            SimpleFeature toWrite = writer.next();
                            toWrite.setAttributes(feature.getAttributes());
                            ((FeatureIdImpl) toWrite.getIdentifier()).setID(feature.getID());
                            toWrite.getUserData().put(Hints.USE_PROVIDED_FID, Boolean.TRUE);
                            toWrite.getUserData().putAll(feature.getUserData());
                            writer.write();
                        }
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
                @Override
                public void close() throws Exception {
                    // 任务停止时关闭连接
                    if (lcLive != null) {
                        lcLive.dispose();
                    }
                }
            });
        env.execute("TransportGeomesaJob");
    }
}
优化建议
  1. 如果需要用事件时间计算窗口,把TumblingProcessingTimeWindows替换为TumblingEventTimeWindows,同时提前配置好水位线生成规则和时间戳提取器即可。
  2. 若2秒窗口数据量过大,单windowAll subtask处理不过来,可以先按业务字段做keyBy,再开并行滚动窗口,多并行度写入Geomesa,吞吐量会更高。
  3. 可以配合添加触发器,设置窗口最大数据条数阈值,超过阈值提前触发写入,避免窗口内数据过多导致OOM。

内容的提问来源于stack exchange,提问作者Ellee

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 13:00:00