Flink实战七之Flink SQL应用
前言

Flink SQL 是 Flink 为简化大数据开发推出的「声明式查询语言」—— 无需编写复杂的 Java/Scala 代码,仅用 SQL 就能完成实时 / 离线数据处理,这也是广告部门降低实时开发成本的核心手段(比如用 Flink SQL 快速实现广告指标统计、归因逻辑)。
一、Flink SQL 核心价值
| 优势 | 广告场景价值 |
|---|---|
| 流批一体 | 一套 SQL 既可以实时统计广告曝光指标,也可以离线计算日账单,无需写两套代码 |
| 低代码 | 运营 / 分析师也能写 SQL 完成广告数据分析,无需依赖开发人员 |
| 高性能 | 底层复用 Flink 引擎的低延迟 / 高吞吐特性,适配广告百万级日志处理 |
| 标准兼容 | 兼容 SQL 标准,支持自定义函数(UDF),适配广告复杂业务逻辑 |
二、Flink SQL 基础用法
1. 环境准备(SQL Client 交互式查询)
Flink 提供「SQL Client」工具,无需写代码,直接在终端执行 SQL:
# 1. 启动 Flink 集群
yangyanping@yangyaningdeAir bin % ./start-cluster.sh
# 2. 启动 Flink SQL Client
yangyanping@yangyaningdeAir bin % ./sql-client.sh embedded
▒▓██▓██▒
▓████▒▒█▓▒▓███▓▒
▓███▓░░ ▒▒▒▓██▒ ▒
░██▒ ▒▒▓▓█▓▓▒░ ▒████
██▒ ░▒▓███▒ ▒█▒█▒
░▓█ ███ ▓░▒██
▓█ ▒▒▒▒▒▓██▓░▒░▓▓█
█░ █ ▒▒░ ███▓▓█ ▒█▒▒▒
████░ ▒▓█▓ ██▒▒▒ ▓███▒
░▒█▓▓██ ▓█▒ ▓█▒▓██▓ ░█░
▓░▒▓████▒ ██ ▒█ █▓░▒█▒░▒█▒
███▓░██▓ ▓█ █ █▓ ▒▓█▓▓█▒
░██▓ ░█░ █ █▒ ▒█████▓▒ ██▓░▒
███░ ░ █░ ▓ ░█ █████▒░░ ░█░▓ ▓░
██▓█ ▒▒▓▒ ▓███████▓░ ▒█▒ ▒▓ ▓██▓
▒██▓ ▓█ █▓█ ░▒█████▓▓▒░ ██▒▒ █ ▒ ▓█▒
▓█▓ ▓█ ██▓ ░▓▓▓▓▓▓▓▒ ▒██▓ ░█▒
▓█ █ ▓███▓▒░ ░▓▓▓███▓ ░▒░ ▓█
██▓ ██▒ ░▒▓▓███▓▓▓▓▓██████▓▒ ▓███ █
▓███▒ ███ ░▓▓▒░░ ░▓████▓░ ░▒▓▒ █▓
█▓▒▒▓▓██ ░▒▒░░░▒▒▒▒▓██▓░ █▓
██ ▓░▒█ ▓▓▓▓▒░░ ▒█▓ ▒▓▓██▓ ▓▒ ▒▒▓
▓█▓ ▓▒█ █▓░ ░▒▓▓██▒ ░▓█▒ ▒▒▒░▒▒▓█████▒
██░ ▓█▒█▒ ▒▓▓▒ ▓█ █░ ░░░░ ░█▒
▓█ ▒█▓ ░ █░ ▒█ █▓
█▓ ██ █░ ▓▓ ▒█▓▓▓▒█░
█▓ ░▓██░ ▓▒ ▓█▓▒░░░▒▓█░ ▒█
██ ▓█▓░ ▒ ░▒█▒██▒ ▓▓
▓█▒ ▒█▓▒░ ▒▒ █▒█▓▒▒░░▒██
░██▒ ▒▓▓▒ ▓██▓▒█▒ ░▓▓▓▓▒█▓
░▓██▒ ▓░ ▒█▓█ ░░▒▒▒
▒▓▓▓▓▓▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒░░▓▓ ▓░▒█░
______ _ _ _ _____ ____ _ _____ _ _ _ BETA
| ____| (_) | | / ____|/ __ \| | / ____| (_) | |
| |__ | |_ _ __ | | __ | (___ | | | | | | | | |_ ___ _ __ | |_
| __| | | | '_ \| |/ / \___ \| | | | | | | | | |/ _ \ '_ \| __|
| | | | | | | | < ____) | |__| | |____ | |____| | | __/ | | | |_
|_| |_|_|_| |_|_|\_\ |_____/ \___\_\______| \_____|_|_|\___|_| |_|\__|
Welcome! Enter 'HELP;' to list all available commands. 'QUIT;' to exit.
Command history file path: /Users/yangyanping/.flink-sql-history
Flink SQL>
2. 核心语法:创建源表(读取数据)+ 结果表(写入数据)
以广告场景为例,先定义「广告点击日志源表(Kafka)」和「指标结果表(MySQL)」,再用 SQL 计算点击率。
步骤 1:创建 Kafka 源表(读取广告点击日志)
SQL语句如下:
-- 修复前(报错):click_time TIMESTAMP(6)
-- 修复后:click_time TIMESTAMP(3)
CREATE TABLE ad_click_log (
ad_id INT,
user_id STRING,
click_time TIMESTAMP(3), -- 核心修改:精度从 6 改为 3(毫秒级)
WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND -- 水位线正常定义
) WITH (
'connector' = 'kafka',
'topic' = 'ad_click_topic',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'ad_click_group',
'scan.startup.mode' = 'latest-offset',
'format' = 'json'
);
-- 定义广告点击日志源表(Kafka 流数据)
Flink SQL> CREATE TABLE ad_click_log (
> ad_id INT,
> user_id STRING,
> click_time TIMESTAMP(3), -- 核心修改:精度从 6 改� 3(毫秒级)
> WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND -- 水位线正常定义
> ) WITH (
> 'connector' = 'kafka',
> 'topic' = 'ad_click_topic',
> 'properties.bootstrap.servers' = 'localhost:9092',
> 'properties.group.id' = 'ad_click_group',
> 'scan.startup.mode' = 'latest-offset',
> 'format' = 'json'
> );
[INFO] Execute statement succeed.
Flink SQL>
步骤 2:创建 MySQL 结果表(写入点击率指标)
SQL语句如下:
-- 定义广告点击率结果表(MySQL 存储)
CREATE TABLE ad_click_rate (
ad_id INT, -- 广告ID
window_time STRING, -- 统计窗口时间
click_count BIGINT, -- 点击量
expose_count BIGINT, -- 曝光量
click_rate DOUBLE, -- 点击率(click_count/expose_count)
PRIMARY KEY (ad_id, window_time) NOT ENFORCED -- 主键(保证幂等写入)
) WITH (
'connector' = 'jdbc', -- 连接器类型
'url' = 'jdbc:mysql://localhost:3306/ad_db',-- MySQL 地址
'table-name' = 'ad_click_rate', -- 目标表名
'username' = 'root', -- 用户名
'password' = 'Yangyanping@1981', -- 密码
'driver' = 'com.mysql.cj.jdbc.Driver' -- 驱动类
);
Flink SQL> CREATE TABLE ad_click_rate (
> ad_id INT, -- 广告ID
> window_time STRING, -- 统计窗口时间
> click_count BIGINT, -- 点击量
> expose_count BIGINT, -- 曝光量
> click_rate DOUBLE, -- 点击率(click_count/expose_count)
> PRIMARY KEY (ad_id, window_time) NOT ENFORCED -- 主键(保证幂等写入)
> ) WITH (
> 'connector' = 'jdbc', -- 连接器类型
> 'url' = 'jdbc:mysql://localhost:3306/ad_db',-- MySQL 地址
> 'table-name' = 'ad_click_rate', -- 目标表名
> 'username' = 'root', -- 用户名
> 'password' = 'Yangyanping@1981', -- 密码
> 'driver' = 'com.mysql.cj.jdbc.Driver' -- 驱动类
> );
[INFO] Execute statement succeed.
Flink SQL>
步骤 3:创建 ad_expose_log 表
参考 ad_click_log 的正确配置,创建曝光日志表(同样修复时间精度问题):
-- 核心:创建广告曝光日志表(Kafka 源表,适配 Flink 水位线精度)
CREATE TABLE ad_expose_log (
ad_id INT, -- 广告ID(和点击表关联)
user_id STRING, -- 曝光用户ID
expose_time TIMESTAMP(3),-- 曝光时间(精度3位,适配水位线)
WATERMARK FOR expose_time AS expose_time - INTERVAL '5' SECOND -- 处理乱序曝光日志
) WITH (
'connector' = 'kafka', -- 连接器:Kafka(广告曝光日志通常存在Kafka)
'topic' = 'ad_expose_topic', -- 替换为你的曝光日志Kafka主题
'properties.bootstrap.servers' = 'localhost:9092', -- Kafka地址
'properties.group.id' = 'ad_expose_group', -- 消费者组(和点击表区分)
'scan.startup.mode' = 'latest-offset', -- 从最新偏移量消费
'format' = 'json' -- 数据格式:JSON(广告日志主流格式)
);
执行后若提示 Create table 'ad_expose_log' successfully.,说明表创建成功。
步骤 4:先确认「所有依赖表」是否都已创建
广告点击率计算需要关联
ad_click_log(点击)和ad_expose_log(曝光)两张表,必须先确保两张表都创建成功:
-- 1. 查看已创建的表(核心校验)
Flink SQL> SHOW TABLES;
- 正常结果(两张表都在):
+---------------+ | table name | +---------------+ | ad_click_log | | ad_expose_log | -- 必须有这行,否则就是没创建 | ad_click_rate | +---------------+ - 异常结果:列表中无
ad_expose_log→ 直接执行下面的建表语句创建。
步骤 5:编写 SQL 计算 5 分钟窗口点击率
MySQL数据库创建表

SQL语句如下:
-- 关联曝光日志表(假设 ad_expose_log 已定义),计算 5 分钟点击率
INSERT INTO ad_click_rate
SELECT
a.ad_id,
COUNT(DISTINCT a.user_id) AS click_count,
COUNT(DISTINCT b.user_id) AS expose_count,
ROUND(COUNT(DISTINCT a.user_id) / COUNT(DISTINCT b.user_id), 4) AS click_rate
FROM ad_click_log a
LEFT JOIN ad_expose_log b
ON a.ad_id = b.ad_id
GROUP BY a.ad_id;
三、广告场景核心实战(Flink SQL 典型用法)
1. 实时归因(最后点击归因)
广告场景核心需求:用户下单前最后一次点击的广告即为归因广告,用 Flink SQL 实现:
-- 定义订单表(Kafka 流)
CREATE TABLE ad_order (
order_id STRING,
user_id STRING,
order_time TIMESTAMP,
WATERMARK FOR order_time AS order_time - INTERVAL '10' SECOND
) WITH ('connector' = 'kafka', ...);
-- 最后点击归因:关联点击表和订单表,取用户下单前最后一次点击的广告
INSERT INTO ad_attribution_result
SELECT
o.order_id,
o.user_id,
LAST_VALUE(c.ad_id) OVER (
PARTITION BY o.user_id
ORDER BY c.click_time
ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
) AS attributed_ad_id,
o.order_time
FROM ad_order o
LEFT JOIN ad_click_log c
ON o.user_id = c.user_id
AND c.click_time <= o.order_time
AND c.click_time >= o.order_time - INTERVAL '1' HOUR; -- 只看下单前1小时内的点击
2. 动态广告规则过滤
通过 Flink SQL 关联 MySQL 中的广告规则表(CDC 同步),过滤无效广告:
-- 定义广告规则表(CDC 实时同步 MySQL 变更)
CREATE TABLE ad_rule (
ad_id INT,
status INT, -- 1=生效,0=失效
bid_price DOUBLE,
PRIMARY KEY (ad_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'localhost',
'port' = '3306',
'username' = 'root',
'password' = '123456',
'database-name' = 'ad_db',
'table-name' = 'ad_rule'
);
-- 过滤失效广告,只统计生效广告的曝光量
SELECT
c.ad_id,
COUNT(*) AS expose_count
FROM ad_expose_log c
JOIN ad_rule r ON c.ad_id = r.ad_id
WHERE r.status = 1 -- 只保留生效广告
GROUP BY c.ad_id;
四、Flink SQL 核心特性(广告场景必知)
1. 窗口函数(统计时间维度指标)
表格
| 窗口类型 | 广告场景用途 |
|---|---|
| 滚动窗口(TUMBLE) | 固定 5/10 分钟统计广告指标(如每 5 分钟点击率) |
| 滑动窗口(HOP) | 每 1 分钟统计过去 5 分钟的广告曝光量 |
| 会话窗口(SESSION) | 统计用户连续点击广告的会话时长(如用户 30 分钟内的点击行为) |
2. 时间语义(处理乱序日志)
- 处理时间:Flink 处理数据的时间(简单但不准确);
- 事件时间:广告日志的实际产生时间(如点击时间),结合水位线(Watermark)处理乱序数据(广告日志跨地域传输必用)。
3. 状态管理(保证计算准确性)
Flink SQL 自动管理状态,支持「精准一次(Exactly-Once)」语义,避免广告指标重复计算 / 丢失。
4. UDF 自定义函数(适配复杂业务)
广告场景中复杂的规则计算(如用户画像标签匹配)可通过 UDF 实现:
// 自定义 UDF:判断用户是否为高价值用户
public class HighValueUserUDF extends ScalarFunction {
public boolean eval(String user_id) {
// 业务逻辑:查询用户消费金额,判断是否为高价值
return checkUserValue(user_id) > 1000;
}
}
在 Flink SQL 中注册并使用:
-- 注册 UDF
CREATE FUNCTION is_high_value_user AS 'com.yyp.HighValueUserUDF';
-- 使用 UDF 筛选高价值用户的广告点击
SELECT ad_id, COUNT(*) FROM ad_click_log
WHERE is_high_value_user(user_id) = TRUE
GROUP BY ad_id;
五、生产环境部署(广告场景)
Flink SQL 可通过「SQL 文件」批量执行,无需嵌入代码:
# 将 SQL 写入文件 ad_sql.sql
# 提交 SQL 作业到 Flink 集群
./bin/flink run -d -c org.apache.flink.table.api.bridge.scala.StreamTableEnvironment ./lib/flink-table-bridge-1.18.1.jar -f ad_sql.sql
总结
Flink SQL 核心要点:
- 核心价值:流批一体、低代码,适配广告场景快速实现实时 / 离线数据处理;
- 基础用法:定义源表(Kafka/CDC)→ 编写计算 SQL → 写入结果表(MySQL/ClickHouse);
- 广告场景核心:窗口统计指标、最后点击归因、动态规则过滤;
- 关键特性:事件时间 + 水位线(处理乱序)、UDF(复杂逻辑)、精准一次(数据准确)。
相比于手写 Java 代码,Flink SQL 能将广告实时开发效率提升 50% 以上,也是广告部门运营 / 开发协作的核心工具。
更多推荐
所有评论(0)