前言

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 核心要点:

  1. 核心价值:流批一体、低代码,适配广告场景快速实现实时 / 离线数据处理;
  2. 基础用法:定义源表(Kafka/CDC)→ 编写计算 SQL → 写入结果表(MySQL/ClickHouse);
  3. 广告场景核心:窗口统计指标、最后点击归因、动态规则过滤;
  4. 关键特性:事件时间 + 水位线(处理乱序)、UDF(复杂逻辑)、精准一次(数据准确)。

相比于手写 Java 代码,Flink SQL 能将广告实时开发效率提升 50% 以上,也是广告部门运营 / 开发协作的核心工具。

更多推荐