概述
FlinkSQL 是 Apache Flink 提供的 SQL 接口,允许用户使用标准 SQL 语法进行流处理和批处理。它降低了 Flink 的使用门槛,让熟悉 SQL 的用户能够快速上手流处理开发。
💡 提示: FlinkSQL 支持 ANSI SQL 标准,同时扩展了流处理相关的语法,如窗口函数、时间函数等。
核心特性
- SQL 标准兼容: 支持 ANSI SQL 标准语法
- 流批一体: 同样的 SQL 可以用于流处理和批处理
- 丰富的连接器: 支持 Kafka、MySQL、Elasticsearch 等多种数据源
- 实时窗口: 支持滚动、滑动、会话等多种窗口类型
- 状态管理: 自动管理查询状态和容错
环境准备
系统要求
- Java: JDK 8 或 11+
- Flink: 1.17.x 版本
- Docker: 用于快速部署(可选)
1. Docker 部署 Flink 集群
创建网络和拉取镜像
docker pull flink:1.17.2
docker network create flink-network
启动 JobManager
docker run \
-itd \
--name=jobmanager \
--publish 8081:8081 \
--network flink-network \
--env FLINK_PROPERTIES="jobmanager.rpc.address: jobmanager" \
flink:1.17.2 jobmanager
启动 TaskManager
docker run \
-itd \
--name=taskmanager \
--network flink-network \
--env FLINK_PROPERTIES="jobmanager.rpc.address: jobmanager" \
flink:1.17.2 taskmanager
启动 SQL Client
docker run \
-it \
--network flink-network \
--env FLINK_PROPERTIES="jobmanager.rpc.address: jobmanager" \
flink:1.17.2 sql-client
配置 Task Slots
默认 Task Slots 为 1,建议调整为 4-8:
# 复制配置文件
docker cp taskmanager:/opt/flink/conf/flink-conf.yaml .
# 修改 taskmanager.numberOfTaskSlots = 5
# 复制回容器
docker cp flink-conf.yaml taskmanager:/opt/flink/conf/
docker restart taskmanager
2. 本地开发环境
Maven 依赖
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>1.17.2</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner</artifactId>
<version>1.17.2</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.2</version>
</dependency>
<!-- Kafka Connector -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-sql-connector-kafka</artifactId>
<version>3.1.0-1.17</version>
</dependency>
</dependencies>
3. 连接器安装
Kafka 连接器
# 下载 Kafka 连接器
wget https://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-kafka/3.1.0-1.17/flink-sql-connector-kafka-3.1.0-1.17.jar
# 复制到 Flink lib 目录(JobManager 和 TaskManager 都需要)
docker cp flink-sql-connector-kafka-3.1.0-1.17.jar jobmanager:/opt/flink/lib/
docker cp flink-sql-connector-kafka-3.1.0-1.17.jar taskmanager:/opt/flink/lib/
# 重启容器
docker restart jobmanager taskmanager
核心概念
1. 表和视图
- Table: FlinkSQL 中的基本数据结构
- View: 虚拟表,基于查询结果创建
- Temporary Table: 临时表,仅在当前会话有效
- Catalog: 元数据管理,存储表、视图等信息
2. 时间语义
- Processing Time: 处理时间,基于系统时钟
- Event Time: 事件时间,基于数据中的时间戳
- Watermark: 水位线,用于处理延迟数据
3. 窗口类型
- Tumbling Window: 滚动窗口,固定大小,无重叠
- Sliding Window: 滑动窗口,固定大小,有重叠
- Session Window: 会话窗口,基于活动间隔
- Over Window: 行级窗口,类似传统数据库的窗口函数
SQL 语法
1. 基本查询
-- 简单查询
SELECT user_id, COUNT(*) as cnt
FROM user_events
GROUP BY user_id;
-- 条件过滤
SELECT *
FROM user_events
WHERE event_type = 'click'
AND event_time > TIMESTAMP '2025-01-01 00:00:00';
-- 排序和限制
SELECT user_id, event_count
FROM user_stats
ORDER BY event_count DESC
LIMIT 10;
2. 窗口查询
-- 滚动窗口(每分钟统计)
SELECT
window_start,
window_end,
COUNT(*) as event_count
FROM TABLE(
TUMBLE(TABLE user_events, DESCRIPTOR(event_time), INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end;
-- 滑动窗口(每30秒统计最近1分钟)
SELECT
window_start,
window_end,
AVG(amount) as avg_amount
FROM TABLE(
HOP(TABLE orders, DESCRIPTOR(order_time), INTERVAL '30' SECOND, INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end;
-- 会话窗口(30秒无活动则结束会话)
SELECT
user_id,
window_start,
window_end,
COUNT(*) as session_events
FROM TABLE(
SESSION(TABLE user_events, DESCRIPTOR(event_time), INTERVAL '30' SECOND)
)
GROUP BY user_id, window_start, window_end;
3. 连接查询
-- 内连接
SELECT u.user_id, u.username, o.order_id, o.amount
FROM users u
JOIN orders o ON u.user_id = o.user_id;
-- 时间窗口连接
SELECT
u.user_id,
u.username,
o.order_id,
o.amount
FROM users u
JOIN orders o ON u.user_id = o.user_id
WHERE u.update_time BETWEEN o.order_time - INTERVAL '10' MINUTE
AND o.order_time + INTERVAL '10' MINUTE;
4. 时间函数
-- 当前时间
SELECT CURRENT_TIMESTAMP, LOCALTIMESTAMP;
-- 时间格式转换
SELECT
DATE_FORMAT(event_time, 'yyyy-MM-dd HH:mm:ss') as formatted_time,
UNIX_TIMESTAMP(event_time) as timestamp_seconds
FROM events;
-- 时间计算
SELECT
event_time,
event_time + INTERVAL '1' HOUR as one_hour_later,
TIMESTAMPDIFF(MINUTE, event_time, CURRENT_TIMESTAMP) as minutes_ago
FROM events;
5. DDL 语句
-- 创建 Catalog(以 Hive 为例,需要 flink-sql-connector-hive 依赖)
CREATE CATALOG hive_catalog WITH (
'type' = 'hive',
'default-database' = 'default',
'hive-conf-dir' = '/opt/hive/conf'
);
USE CATALOG hive_catalog;
-- 创建数据库
CREATE DATABASE IF NOT EXISTS realtime COMMENT '实时数仓';
USE realtime;
-- 注册 UDF(JAR 需要先放入 lib 目录,或用 USING JAR 指定)
CREATE TEMPORARY FUNCTION mask_phone AS 'com.example.udf.MaskPhone'
LANGUAGE JAVA USING JAR 'file:///opt/udf/my-udf.jar';
-- 删除对象
DROP TABLE IF EXISTS user_events;
DROP VIEW IF EXISTS active_users;
DROP TEMPORARY FUNCTION IF EXISTS mask_phone;
6. DML 语句
-- 插入常量数据
INSERT INTO user_stats VALUES (1, 10, TIMESTAMP '2024-01-01 00:00:00');
-- 插入查询结果(流作业会持续运行)
INSERT INTO user_stats
SELECT user_id, COUNT(*), MAX(event_time)
FROM user_events
GROUP BY user_id;
-- 写入分区表的指定分区
INSERT INTO dws_orders PARTITION (dt = '2024-01-01')
SELECT order_id, amount FROM orders;
-- 批模式下覆盖写入(仅 Filesystem/Hive 等连接器支持)
INSERT OVERWRITE dws_orders PARTITION (dt = '2024-01-01')
SELECT order_id, amount FROM orders;
-- 多个 INSERT 合并成一个作业,共享同一份源数据
EXECUTE STATEMENT SET
BEGIN
INSERT INTO sink_a SELECT * FROM user_events WHERE event_type = 'click';
INSERT INTO sink_b SELECT * FROM user_events WHERE event_type = 'buy';
END;
UPDATE 和 DELETE 只在批模式(SET 'execution.runtime-mode' = 'batch';)下可用,并且要求目标连接器实现了行级修改能力,流作业中不能使用。
7. 常用查询模式
Top-N
流上的 Top-N 必须用 ROW_NUMBER() 加外层过滤,单纯的 ORDER BY ... LIMIT 在流模式下只能按时间属性排序。
SELECT category, product_id, sales
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY category ORDER BY sales DESC) AS rn
FROM product_sales
)
WHERE rn <= 3;
Window Top-N
先用窗口 TVF 聚合,再在每个窗口内排名:
SELECT *
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY cnt DESC) AS rn
FROM (
SELECT window_start, window_end, product_id, COUNT(*) AS cnt
FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '1' HOUR))
GROUP BY window_start, window_end, product_id
)
)
WHERE rn <= 3;
去重(Deduplication)
按主键保留第一条(或最后一条)记录。ORDER BY 必须是时间属性字段,ASC 保留第一条,DESC 保留最后一条。
SELECT order_id, user_id, amount, proc_time
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY proc_time ASC) AS rn
FROM orders
)
WHERE rn = 1;
集合运算
SELECT user_id FROM orders_a
UNION -- 去重合并,UNION ALL 保留重复
SELECT user_id FROM orders_b;
SELECT user_id FROM orders_a
INTERSECT -- 交集
SELECT user_id FROM orders_b;
SELECT user_id FROM orders_a
EXCEPT -- 差集
SELECT user_id FROM orders_b;
模式识别(MATCH_RECOGNIZE)
检测同一股票价格连续上涨的事件序列:
SELECT *
FROM ticker
MATCH_RECOGNIZE (
PARTITION BY symbol
ORDER BY rowtime
MEASURES
START_ROW.rowtime AS start_time,
LAST(PRICE_UP.rowtime) AS end_time,
LAST(PRICE_UP.price) AS top_price
ONE ROW PER MATCH
AFTER MATCH SKIP PAST LAST ROW
PATTERN (START_ROW PRICE_UP+)
DEFINE
PRICE_UP AS PRICE_UP.price > PREV(PRICE_UP.price)
);
连接器配置
1. Kafka 连接器
-- 创建 Kafka 源表
CREATE TABLE user_events (
user_id BIGINT,
event_type STRING,
event_time TIMESTAMP(3),
properties MAP<STRING, STRING>,
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user_events',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink_consumer_group',
'scan.startup.mode' = 'earliest-offset',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
-- 创建 Kafka 结果表
CREATE TABLE user_stats (
user_id BIGINT,
event_count BIGINT,
last_event_time TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'topic' = 'user_stats',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
启动位置 scan.startup.mode
| 取值 | 含义 | 适用场景 |
|---|---|---|
group-offsets(默认) | 从消费者组已提交的偏移量继续消费,分区没有提交记录时按 properties.auto.offset.reset 处理 | 作业重启后从中断位置继续 |
earliest-offset | 从每个分区最早的偏移量开始 | 回放全部历史数据 |
latest-offset | 从作业启动时的最新偏移量开始 | 只处理新数据 |
timestamp | 从指定时间戳之后的第一条消息开始,配合 scan.startup.timestamp-millis | 从某个时间点重放 |
specific-offsets | 为每个分区指定起始偏移量,配合 scan.startup.specific-offsets | 精确控制消费位置 |
-- 从指定时间点开始
'scan.startup.mode' = 'timestamp',
'scan.startup.timestamp-millis' = '1704038400000'
-- 为每个分区指定偏移量
'scan.startup.mode' = 'specific-offsets',
'scan.startup.specific-offsets' = 'partition:0,offset:42;partition:1,offset:300'
元数据列与嵌套 JSON
CREATE TABLE kafka_source (
count_num INT,
details ROW<`user` STRING, status STRING>, -- 嵌套 JSON 对象映射为 ROW
kafka_ts TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' VIRTUAL, -- Kafka 消息时间戳
kafka_partition INT METADATA FROM 'partition' VIRTUAL
) WITH (
'connector' = 'kafka',
'topic' = 'kafka_demo',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'format' = 'json'
);
-- 访问嵌套字段
SELECT count_num, details.`user`, details.status, kafka_ts FROM kafka_source;
2. MySQL 连接器
-- 创建 MySQL 维表
CREATE TABLE user_info (
user_id BIGINT,
username STRING,
email STRING,
create_time TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/test',
'table-name' = 'user_info',
'username' = 'root',
'password' = 'password',
'lookup.cache.max-rows' = '1000',
'lookup.cache.ttl' = '1h'
);
-- 创建 MySQL 结果表
CREATE TABLE daily_stats (
stat_date DATE,
user_count BIGINT,
event_count BIGINT,
PRIMARY KEY (stat_date) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/test',
'table-name' = 'daily_stats',
'username' = 'root',
'password' = 'password'
);
3. Elasticsearch 连接器
-- 创建 Elasticsearch 结果表
CREATE TABLE user_behavior_es (
user_id BIGINT,
behavior_type STRING,
event_time TIMESTAMP(3),
properties MAP<STRING, STRING>
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://localhost:9200',
'index' = 'user_behavior',
'document-type' = '_doc',
'bulk-flush.max-actions' = '1000',
'bulk-flush.max-size' = '2mb',
'bulk-flush.interval' = '1s'
);
4. Filesystem 连接器
-- 读写本地或 HDFS 上的文件,支持 csv、json、parquet、orc 等格式
CREATE TABLE file_orders (
user_id STRING,
order_amount DOUBLE,
dt STRING
) PARTITIONED BY (dt) WITH (
'connector' = 'filesystem',
'path' = 'file:///data/orders', -- 或 hdfs://namenode:8020/data/orders
'format' = 'csv'
);
实战示例
1. 实时用户行为分析
场景描述
分析用户实时行为,计算每分钟的 PV、UV,并输出到 Kafka 和 Elasticsearch。
数据生产者(Python)
import json
import time
import random
from datetime import datetime
from kafka import KafkaProducer
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
def generate_user_event():
return {
'user_id': random.randint(1, 1000),
'event_type': random.choice(['click', 'view', 'purchase']),
'page_id': random.randint(1, 100),
'event_time': datetime.now().isoformat(),
'properties': {
'device': random.choice(['mobile', 'desktop']),
'browser': random.choice(['chrome', 'firefox', 'safari'])
}
}
# 持续发送数据
while True:
event = generate_user_event()
producer.send('user_events', value=event)
time.sleep(0.1)
FlinkSQL 处理逻辑
-- 创建源表
CREATE TABLE user_events (
user_id BIGINT,
event_type STRING,
page_id BIGINT,
event_time TIMESTAMP(3),
properties MAP<STRING, STRING>,
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user_events',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json',
'json.timestamp-format.standard' = 'ISO-8601'
);
-- 创建结果表
CREATE TABLE pv_uv_stats (
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
pv BIGINT,
uv BIGINT,
event_type STRING
) WITH (
'connector' = 'kafka',
'topic' = 'pv_uv_stats',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
-- 实时 PV/UV 统计
INSERT INTO pv_uv_stats
SELECT
window_start,
window_end,
COUNT(*) as pv,
COUNT(DISTINCT user_id) as uv,
event_type
FROM TABLE(
TUMBLE(TABLE user_events, DESCRIPTOR(event_time), INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end, event_type;
2. 实时订单监控
场景描述
监控订单流,计算每5分钟的订单金额、笔数,并检测异常订单。
-- 订单流表
CREATE TABLE orders (
order_id STRING,
user_id BIGINT,
amount DECIMAL(10,2),
order_time TIMESTAMP(3),
status STRING,
WATERMARK FOR order_time AS order_time - INTERVAL '10' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
-- 用户信息维表
CREATE TABLE user_profiles (
user_id BIGINT,
username STRING,
user_level STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/ecommerce',
'table-name' = 'user_profiles'
);
-- 订单统计
CREATE VIEW order_stats AS
SELECT
window_start,
window_end,
COUNT(*) as order_count,
SUM(amount) as total_amount,
AVG(amount) as avg_amount,
MAX(amount) as max_amount
FROM TABLE(
TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '5' MINUTE)
)
WHERE status = 'completed'
GROUP BY window_start, window_end;
-- 异常订单检测(大额订单)
CREATE VIEW abnormal_orders AS
SELECT
o.order_id,
o.user_id,
u.username,
o.amount,
o.order_time
FROM orders o
LEFT JOIN user_profiles FOR SYSTEM_TIME AS OF o.order_time AS u
ON o.user_id = u.user_id
WHERE o.amount > 10000
AND u.user_level != 'VIP';
3. 实时数据清洗
场景描述
清洗原始日志数据,过滤无效数据,格式化后输出到数据湖。
-- 原始日志表
CREATE TABLE raw_logs (
log_time TIMESTAMP(3),
level STRING,
message STRING,
source STRING,
properties MAP<STRING, STRING>,
WATERMARK FOR log_time AS log_time - INTERVAL '1' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'raw_logs',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
-- 清洗后的日志表
CREATE TABLE cleaned_logs (
log_time TIMESTAMP(3),
level STRING,
category STRING,
error_code STRING,
user_id BIGINT,
session_id STRING
) WITH (
'connector' = 'iceberg',
'catalog-name' = 'hadoop_catalog',
'database-name' = 'logs',
'table-name' = 'cleaned_logs'
);
-- 数据清洗逻辑
INSERT INTO cleaned_logs
SELECT
log_time,
level,
CASE
WHEN message LIKE '%user%' THEN 'user_action'
WHEN message LIKE '%system%' THEN 'system'
WHEN message LIKE '%error%' THEN 'error'
ELSE 'unknown'
END as category,
REGEXP_EXTRACT(message, 'error_code:(\d+)', 1) as error_code,
CAST(properties['user_id'] AS BIGINT) as user_id,
properties['session_id'] as session_id
FROM raw_logs
WHERE level IN ('INFO', 'WARN', 'ERROR')
AND message IS NOT NULL
AND char_length(message) > 10;
性能优化
1. 查询优化
-- 使用局部聚合优化
SET 'table.optimizer.agg-phase-strategy' = 'TWO_PHASE';
-- 启用MiniBatch聚合
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '1s';
SET 'table.exec.mini-batch.size' = '1000';
-- 状态清理
SET 'table.exec.state.ttl' = '1h';
2. 并行度配置
-- 设置全局并行度
SET 'parallelism.default' = '4';
-- 针对特定算子设置并行度
CREATE TABLE high_throughput_table (...)
WITH (..., 'sink.parallelism' = '8');
3. 内存管理
-- 配置状态后端
SET 'state.backend' = 'rocksdb';
SET 'state.checkpoints.dir' = 'hdfs://namenode:port/checkpoints';
-- 配置TaskManager内存
SET 'taskmanager.memory.process.size' = '2g';
SET 'taskmanager.memory.managed.fraction' = '0.4';
4. Checkpoint 优化
-- Checkpoint配置
SET 'execution.checkpointing.interval' = '10s';
SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
SET 'execution.checkpointing.timeout' = '60s';
SET 'execution.checkpointing.max-concurrent-checkpoints' = '1';
最佳实践
1. 表设计原则
合理使用主键
-- 带主键的表支持 Upsert 语义
CREATE TABLE user_stats (
user_id BIGINT,
total_orders BIGINT,
last_order_time TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (...);
选择合适的时间语义
-- Event Time(推荐用于业务时间)
CREATE TABLE events (
event_time TIMESTAMP(3),
...,
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (...);
-- Processing Time(推荐用于监控场景)
CREATE TABLE monitoring (
proc_time AS PROCTIME(),
...
) WITH (...);
2. 查询设计模式
流式 ETL
-- 数据转换和清洗
INSERT INTO clean_data
SELECT
user_id,
UPPER(username) as username,
CASE WHEN age < 18 THEN 'minor' ELSE 'adult' END as age_group,
event_time
FROM raw_data
WHERE user_id IS NOT NULL;
实时报表
-- 多维度实时统计
CREATE VIEW real_time_dashboard AS
SELECT
DATE_FORMAT(window_start, 'yyyy-MM-dd HH:mm') as time_window,
event_type,
COUNT(*) as event_count,
COUNT(DISTINCT user_id) as unique_users,
AVG(CAST(properties['duration'] AS DOUBLE)) as avg_duration
FROM TABLE(
TUMBLE(TABLE user_events, DESCRIPTOR(event_time), INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end, event_type;
流式机器学习特征
-- 用户行为特征提取
CREATE VIEW user_features AS
SELECT
user_id,
COUNT(*) OVER (
PARTITION BY user_id
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW
) as events_last_hour,
COUNT(DISTINCT page_id) OVER (
PARTITION BY user_id
ORDER BY event_time
RANGE BETWEEN INTERVAL '1' DAY PRECEDING AND CURRENT ROW
) as unique_pages_last_day
FROM user_events;
3. 错误处理
数据质量检查
-- 创建侧输出表处理错误数据
CREATE TABLE error_data (
raw_data STRING,
error_message STRING,
error_time TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'topic' = 'error_data',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
-- 在查询中使用 TRY 函数处理异常
SELECT
user_id,
TRY_CAST(properties['age'] AS INT) as age,
event_time
FROM user_events
WHERE TRY_CAST(properties['age'] AS INT) IS NOT NULL;
4. 运维监控
延迟监控
-- 创建延迟监控表
CREATE TABLE latency_monitor (
checkpoint_time TIMESTAMP(3),
source_watermark TIMESTAMP(3),
processing_time TIMESTAMP(3),
latency_seconds BIGINT
) WITH (...);
-- 计算处理延迟
INSERT INTO latency_monitor
SELECT
CURRENT_TIMESTAMP as checkpoint_time,
CURRENT_WATERMARK(event_time) as source_watermark,
PROCTIME() as processing_time,
TIMESTAMPDIFF(SECOND, CURRENT_WATERMARK(event_time), PROCTIME()) as latency_seconds
FROM user_events;
常见问题
问题1:时间语义选择
❓ 问题: 什么时候使用 Event Time vs Processing Time?
解答:
- Event Time: 业务数据分析,需要准确的业务时间
- Processing Time: 系统监控,对延迟敏感的场景
问题2:水位线设置
❓ 问题: 如何设置合适的水位线延迟?
解答:
-- 根据数据延迟情况设置
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND -- 5秒延迟
WATERMARK FOR event_time AS event_time - INTERVAL '1' MINUTE -- 1分钟延迟
问题3:状态过大
❓ 问题: 状态不断增长怎么办?
解答:
-- 设置状态 TTL
SET 'table.exec.state.ttl' = '2h';
-- 使用窗口聚合替代无界聚合
SELECT window_start, window_end, COUNT(*)
FROM TABLE(TUMBLE(...))
GROUP BY window_start, window_end;
问题4:连接器配置
❓ 问题: Kafka 连接器的常见配置问题?
解答:
-- 确保正确的格式配置
'format' = 'json',
'json.fail-on-missing-field' = 'false',
'json.ignore-parse-errors' = 'true'
-- 设置合适的消费组
'properties.group.id' = 'unique_consumer_group',
'scan.startup.mode' = 'earliest-offset' -- 或 'latest-offset'
相关文章
系列导航:参见 Flink 系列导航
Flink 核心文档
- Flink 工作流程剖析 - Flink 底层原理和架构
- Flink DataStream API - 流处理编程 API
- FlinkSQL 内置函数 - SQL 函数完整指南
Flink 高级功能
- Flink Table & SQL API 实时数仓 - 数仓建设方案
- Flink CDC - 变更数据捕获
相关技术
- Kafka 简明教程 - 消息队列
- MySQL 使用指南 - 关系数据库
- Elasticsearch + Kibana - 搜索引擎
总结
FlinkSQL 为流处理提供了简洁的 SQL 接口,核心要素包括:
- 标准 SQL: 降低学习门槛,快速上手
- 丰富连接器: 支持多种数据源和目标系统
- 实时窗口: 提供灵活的时间窗口聚合
- 状态管理: 自动管理查询状态和容错
- 性能优化: 支持多种优化策略
掌握 FlinkSQL 能够快速构建实时数据处理管道,满足从简单的数据清洗到复杂的实时分析等各种需求。