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"); } }
优化建议
- 如果需要用事件时间计算窗口,把
TumblingProcessingTimeWindows替换为TumblingEventTimeWindows,同时提前配置好水位线生成规则和时间戳提取器即可。 - 若2秒窗口数据量过大,单windowAll subtask处理不过来,可以先按业务字段做keyBy,再开并行滚动窗口,多并行度写入Geomesa,吞吐量会更高。
- 可以配合添加触发器,设置窗口最大数据条数阈值,超过阈值提前触发写入,避免窗口内数据过多导致OOM。
内容的提问来源于stack exchange,提问作者Ellee
相关产品推荐
相关产品推荐

