更新时间:2026-08-21 GMT+08:00
分享

Flink SQL语法增强

场景描述

FlinkSQL是开发实时流处理作业的主要方式,但标准FlinkSQL在数据分区控制、迟到数据处理和Source并发调节方面存在一些限制:

  • 数据分区控制:默认的Forward分区策略可能导致下游算子的数据倾斜和背压。
  • 迟到数据处理:在窗口计算中,EventTime落后于Watermark的迟到数据默认被丢弃,无法单独处理。
  • Source并发调节:Source算子的并发数无法独立调整,与下游算子并发不匹配时会引起性能问题。

FlinkSQL提供了DISTRIBUTEBY分区Hint、窗口函数迟到数据处理和Source并发设置SQL语法增强特性,解决数据倾斜、背压和数据丢失等常见问题。

基本概念

  • DISTRIBUTEBY:FlinkSQL Hint特性,通过在SELECT语句中添加 /*+ DISTRIBUTEBY('field') */ 指定分区字段,Flink按该字段Hash值将数据分发到下游算子。支持单字段和多字段分区,适用于仅需数据分区无需聚合的场景。
  • 迟到数据:在事件时间语义下,数据的事件时间(EventTime)落后于当前Watermark的数据称为迟到数据。迟到数据无法参与其所属窗口的计算,标准FlinkSQL中默认丢弃。
  • 数据倾斜:数据分布不均匀导致部分下游算子Task负载过重而其他Task空闲,影响作业整体性能。
  • 背压:下游算子处理速度跟不上上游数据生产速度时产生的反向压力,导致上游算子阻塞,影响作业吞吐。

约束与限制

本章节适用于MRS 3.3.0-LTS及之后版本。

FlinkSQL DISTRIBUTEBY

FlinkSQL新增DISTRIBUTEBY特性,根据指定的字段进行分区,支持单字段及多字段,解决数据仅需要分区的场景。示例如下:

SELECT /*+ DISTRIBUTEBY('id') */ id, name FROM t1;
SELECT /*+ DISTRIBUTEBY('id', 'name') */ id, name FROM t1;
SELECT /*+ DISTRIBUTEBY('id1') */ id as id1, name FROM t1;

FlinkSQL窗口函数支持迟到数据

FlinkSQL新增窗口函数支持迟到数据特性,解决迟到数据需要处理的场景。目前支持TUMBLE、HOP、OVER、CUMULATE窗口函数的迟到数据,示例如下:

CREATE TABLE T1 (
 `int` INT,
 `double` DOUBLE,
 `float` FLOAT,
 `bigdec` DECIMAL(10, 2),
 `string` STRING,
 `name` STRING,
 `rowtime` TIMESTAMP(3),
 WATERMARK for `rowtime` AS `rowtime` - INTERVAL '1' SECOND
) WITH ( 
 'connector' = 'values',
);

-- 该Sink的字段必须和窗口的输入数据保持一致,但顺序不要求一致
CREATE TABLE LD_SINK(
 `float` FLOAT, `string` STRING, `name` STRING,  `rowtime` TIMESTAMP(3)
) WITH ( 
 'connector' = 'print',
);

SELECT  /*+ LATE_DATA_SINK('sink.name'='LD_SINK') */
  `name`,
  MIN(`float`),
  COUNT(DISTINCT `string`)
FROM TABLE(
  TUMBLE(TABLE T1, DESCRIPTOR(rowtime), INTERVAL '5' SECOND))
GROUP BY `name`, window_start, window_end

该特性还支持窗口接收到迟到数据时输出当前窗口的开始时间和结束时间,可通过添加在Hint中'window.start.field'和'window.end.field'使用,字段类型必须是timestamp,示例如下:

CREATE TABLE LD_SINK(
 `float` FLOAT, `string` STRING, `name` STRING,  `rowtime` TIMESTAMP(3), `windowStart` TIMESTAMP(3), `windowEnd` TIMESTAMP(3)
) WITH ( 
 'connector' = 'print',
);

SELECT  /*+ LATE_DATA_SINK('sink.name'='LD_SINK', 'window.start.field'='windowStart', 'window.end.field'='windowEnd') */
  `name`,
  MIN(`float`),
  COUNT(DISTINCT `string`)
FROM TABLE(
  TUMBLE(TABLE T1, DESCRIPTOR(rowtime), INTERVAL '5' SECOND))
GROUP BY `name`, window_start, window_end

FlinkSQL支持设置Source的并发

本章节适用于MRS 3.3.0-LTS及之后版本。

FlinkSQL支持通过使用参数“source.parallelism”设置Source算子的并发数,解决下游算子的并发数引起的一些问题,例如下游算子发送数据倾斜、背压、作业性能慢等问题。

该特性会将Source和下游算子的Forward分区改为Rebalance分区,所以当Source算子的并发数和下游算子的并发数(parallelism数)不一致时,且作业不允许数据乱序,需要在启用该特性的同时开启DISTRIBUTEBY特性,可参考Flink SQL语法增强

如设置Source并发数为“2”并开启DISTRIBUTEBY特性:

CREATE TABLE KafkaSource (
`user_id` VARCHAR,
`user_name` VARCHAR,
 `age` INT
 ) WITH ( 
 'connector' = 'kafka',  
 'topic' = 'test_source', 
 'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',  
 'properties.group.id' = 'testGroup',  
 'scan.startup.mode' = 'latest-offset',  
 'format' = 'csv',  
 'properties.sasl.kerberos.service.name' = 'kafka', 
 'properties.security.protocol' = 'SASL_PLAINTEXT', 
 'properties.kerberos.domain.name' = 'hadoop.系统域名',
 -- 设置Source并发数
 'source.parallelism' = '2'
 ); 
CREATE TABLE KafkaSink( 
  `user_id` VARCHAR, 
  `user_name` VARCHAR,  
 `age` INT
 ) WITH ( 
  'connector' = 'kafka', 
  'topic' = 'test_sink', 
  'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号', 
  'value.format' = 'csv', 
  'properties.sasl.kerberos.service.name' = 'kafka',
   'properties.security.protocol' = 'SASL_PLAINTEXT',
   'properties.kerberos.domain.name' = 'hadoop.系统域名'
 ); 
-- Insert into KafkaSink select user_id, user_name, age from KafkaSource;(未开启DISTRIBUTEBY特性-- 开启DISTRIBUTEBY特性
Insert into KafkaSink select/*+ DISTRIBUTEBY('user_id') */ user_id, user_name, age from KafkaSource;
  • Kafka Broker实例业务IP地址及端口号说明:
    • 服务的实例IP地址可通过登录FusionInsight Manager后,单击“集群 > 服务 > Kafka > 实例”,在实例列表页面中查询。
    • 集群已启用Kerberos认证(安全模式)时Broker端口为“sasl.port”参数的值,默认为“21007”。
    • 集群未启用Kerberos认证(普通模式)时Broker端口为“port”的值,默认为“9092”。如果配置端口号为9092,则需要配置“allow.everyone.if.no.acl.found”参数为true,具体操作如下:

      登录FusionInsight Manager系统,选择“集群 > 服务 > Kafka > 配置 > 全部配置”,搜索“allow.everyone.if.no.acl.found”配置,修改参数值为true,保存配置即可。

  • 系统域名:可登录FusionInsight Manager,选择“系统 > 权限 > 域和互信”,查看“本端域”参数,即为当前系统域名。

常见问题

Q1:使用DISTRIBUTEBY后数据分布发生了什么变化?

A:DISTRIBUTEBY将数据从默认的Forward分区(Source与下游1对1直连)改为按指定字段Hash分区。相同字段值的数据会被路由到同一下游Task,不同字段值的数据会均匀分布到各下游Task,从而解决数据倾斜问题。

Q2:迟到数据的判定标准是什么?

A:Flink通过Watermark机制判定迟到数据。Watermark = 最大已观察EventTime - 允许延迟时间。当新到达数据的事件时间(EventTime)小于当前Watermark时,该数据被判定为迟到数据。例如,Watermark为10:00:05,则事件时间早于10:00:05的数据为迟到数据。

Q3:设置Source并发数后为什么需要同时开启DISTRIBUTEBY?

A:设置source.parallelism后,Source与下游算子的分区策略从Forward(1对1直连)变为Rebalance(轮询分发),Rebalance不保证数据顺序。当Source并发数与下游并发数不一致且作业不允许数据乱序时,需同时开启DISTRIBUTEBY按指定字段Hash分区,保证相同Key的数据路由到同一下游Task,避免数据乱序。

相关文档