全部笔记All notes

FlinkSQL 简明教程

阅读 11m 49s11m 49s read

概述

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: 用于快速部署(可选)

创建网络和拉取镜像

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 系列导航

相关技术


总结

FlinkSQL 为流处理提供了简洁的 SQL 接口,核心要素包括:

  • 标准 SQL: 降低学习门槛,快速上手
  • 丰富连接器: 支持多种数据源和目标系统
  • 实时窗口: 提供灵活的时间窗口聚合
  • 状态管理: 自动管理查询状态和容错
  • 性能优化: 支持多种优化策略

掌握 FlinkSQL 能够快速构建实时数据处理管道,满足从简单的数据清洗到复杂的实时分析等各种需求。