Flink Table & SQL API 实时数仓完整指南
概述
Flink Table & SQL API 提供了声明式的流批统一编程接口,让开发者能够使用熟悉的 SQL 语法进行实时和批量数据处理。它广泛应用于实时数据仓库、实时 OLAP、流式 ETL 等场景。
💡 提示: Table API 和 SQL 在 Flink 中是等价的,SQL 查询会被转换为 Table API 程序,最终编译为 DataStream 程序执行。
核心特性
- 流批统一: 同样的 SQL 可以处理无界流和有界批数据
- 标准 SQL: 支持 ANSI SQL 标准,包括 DDL、DML、DQL
- 丰富的内置函数: 时间函数、窗口函数、聚合函数等
- 扩展性强: 支持自定义函数、连接器、格式
- 性能优化: 查询优化器、状态管理、增量计算
Table & SQL API 基础
核心概念
1. TableEnvironment
管理表的上下文环境,负责:
- 注册和管理表
- 执行 SQL 查询
- 配置优化参数
- 管理 UDF 函数
2. Table
表示表格化的数据,可以:
- 从 DataStream 转换而来
- 通过 SQL DDL 创建
- 从外部系统读取
3. Catalog
元数据管理,包括:
- 数据库
- 表
- 视图
- 函数
环境配置
创建 TableEnvironment
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
// 方式1: 从 StreamExecutionEnvironment 创建
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 方式2: 使用 EnvironmentSettings
EnvironmentSettings settings = EnvironmentSettings
.newInstance()
.inStreamingMode() // 或 .inBatchMode()
.build();
TableEnvironment tableEnv = TableEnvironment.create(settings);
配置优化参数
import org.apache.flink.configuration.Configuration;
Configuration configuration = tableEnv.getConfig().getConfiguration();
// 设置并行度
configuration.setString("parallelism.default", "4");
// 启用 MiniBatch 优化
configuration.setString("table.exec.mini-batch.enabled", "true");
configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
configuration.setString("table.exec.mini-batch.size", "5000");
// 设置空闲状态保留时间
configuration.setString("table.exec.state.ttl", "1 hour");
创建和管理表
1. 从 DataStream 创建表
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.Schema;
import static org.apache.flink.table.api.Expressions.$;
// 定义 POJO
public class Order {
public String orderId;
public String userId;
public Double amount;
public Long orderTime;
// constructors, getters, setters...
}
// 创建 DataStream
DataStream<Order> orderStream = env.fromElements(
new Order("order1", "user1", 100.0, System.currentTimeMillis()),
new Order("order2", "user2", 200.0, System.currentTimeMillis())
);
// 方式1: 自动推断 schema
Table orderTable = tableEnv.fromDataStream(orderStream);
// 方式2: 指定字段和时间属性
Table orderTableWithTime = tableEnv.fromDataStream(
orderStream,
Schema.newBuilder()
.column("orderId", DataTypes.STRING())
.column("userId", DataTypes.STRING())
.column("amount", DataTypes.DOUBLE())
.column("orderTime", DataTypes.BIGINT())
.columnByExpression("proctime", "PROCTIME()") // 处理时间
.columnByExpression("rowtime", "TO_TIMESTAMP_LTZ(orderTime, 3)") // 事件时间
.watermark("rowtime", "rowtime - INTERVAL '5' SECOND") // 水印
.build()
);
// 注册为临时视图
tableEnv.createTemporaryView("orders", orderTableWithTime);
2. 使用 DDL 创建表
-- 创建 Kafka 源表
CREATE TABLE kafka_orders (
order_id STRING,
user_id STRING,
amount DOUBLE,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink-consumer',
'format' = 'json',
'json.fail-on-missing-field' = 'false',
'json.ignore-parse-errors' = 'true'
);
-- 创建 MySQL 维表
CREATE TABLE mysql_users (
user_id STRING,
user_name STRING,
user_level INT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/test',
'table-name' = 'users',
'username' = 'root',
'password' = 'password',
'lookup.cache.max-rows' = '1000',
'lookup.cache.ttl' = '10 min'
);
-- 创建 Elasticsearch 结果表
CREATE TABLE es_order_stats (
user_id STRING,
order_count BIGINT,
total_amount DOUBLE,
stats_time TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://localhost:9200',
'index' = 'order_stats'
);
SQL 查询基础
基本查询
// SELECT 查询
Table result = tableEnv.sqlQuery(
"SELECT user_id, COUNT(*) as order_count, SUM(amount) as total_amount " +
"FROM orders " +
"GROUP BY user_id"
);
// 执行 INSERT
tableEnv.executeSql(
"INSERT INTO es_order_stats " +
"SELECT user_id, COUNT(*) as order_count, SUM(amount) as total_amount, CURRENT_TIMESTAMP " +
"FROM orders " +
"GROUP BY user_id"
);
// 将结果转换回 DataStream
DataStream<Row> resultStream = tableEnv.toDataStream(result);
高级 SQL 特性
窗口函数
1. 滚动窗口 (TUMBLE)
-- 每小时统计
SELECT
user_id,
TUMBLE_START(rowtime, INTERVAL '1' HOUR) as window_start,
TUMBLE_END(rowtime, INTERVAL '1' HOUR) as window_end,
COUNT(*) as order_count,
SUM(amount) as total_amount
FROM orders
GROUP BY
user_id,
TUMBLE(rowtime, INTERVAL '1' HOUR);
2. 滑动窗口 (HOP)
-- 每30分钟统计过去1小时
SELECT
user_id,
HOP_START(rowtime, INTERVAL '30' MINUTE, INTERVAL '1' HOUR) as window_start,
HOP_END(rowtime, INTERVAL '30' MINUTE, INTERVAL '1' HOUR) as window_end,
COUNT(*) as order_count,
AVG(amount) as avg_amount
FROM orders
GROUP BY
user_id,
HOP(rowtime, INTERVAL '30' MINUTE, INTERVAL '1' HOUR);
3. 会话窗口 (SESSION)
-- 会话间隔30分钟
SELECT
user_id,
SESSION_START(rowtime, INTERVAL '30' MINUTE) as session_start,
SESSION_END(rowtime, INTERVAL '30' MINUTE) as session_end,
COUNT(*) as event_count
FROM user_events
GROUP BY
user_id,
SESSION(rowtime, INTERVAL '30' MINUTE);
时间属性
处理时间 vs 事件时间
-- 处理时间
CREATE TABLE process_time_table (
id INT,
data STRING,
proc_time AS PROCTIME() -- 处理时间
) WITH (...);
-- 事件时间
CREATE TABLE event_time_table (
id INT,
data STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND -- 水印
) WITH (...);
复杂查询
1. JOIN 操作
-- 常规 JOIN
SELECT
o.order_id,
o.user_id,
u.user_name,
o.amount
FROM orders o
JOIN users u ON o.user_id = u.user_id;
-- 时态表 JOIN (维表关联)
SELECT
o.order_id,
o.user_id,
u.user_name,
u.user_level,
o.amount
FROM orders AS o
JOIN mysql_users FOR SYSTEM_TIME AS OF o.order_time AS u
ON o.user_id = u.user_id;
-- 区间 JOIN
SELECT
o.order_id,
p.payment_id,
o.amount
FROM orders o, payments p
WHERE o.order_id = p.order_id
AND p.payment_time BETWEEN o.order_time AND o.order_time + INTERVAL '1' HOUR;
2. 窗口函数 (Over)
-- 排名函数
SELECT
user_id,
order_time,
amount,
ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY amount DESC) as rank
FROM orders;
-- 累计统计
SELECT
user_id,
order_time,
amount,
SUM(amount) OVER (
PARTITION BY user_id
ORDER BY order_time
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
) as cumulative_amount
FROM orders;
自定义函数扩展
标量函数 (UDF)
标量函数接收零个或多个参数,返回单个值。
import org.apache.flink.table.functions.ScalarFunction;
import org.apache.flink.table.annotation.DataTypeHint;
import org.apache.flink.table.annotation.FunctionHint;
@FunctionHint(
input = {@DataTypeHint("STRING"), @DataTypeHint("INT")},
output = @DataTypeHint("STRING")
)
public class HashFunction extends ScalarFunction {
public String eval(String input, Integer salt) {
if (input == null) {
return null;
}
return DigestUtils.sha256Hex(input + salt);
}
// 支持重载
public String eval(String input) {
return eval(input, 0);
}
}
// 注册函数
tableEnv.createTemporarySystemFunction("hash", HashFunction.class);
// 使用函数
tableEnv.sqlQuery("SELECT user_id, hash(user_id, 123) as hashed_id FROM users");
表值函数 (UDTF)
表值函数可以返回多行数据。
import org.apache.flink.table.functions.TableFunction;
import org.apache.flink.types.Row;
@FunctionHint(output = @DataTypeHint("ROW<word STRING, length INT>"))
public class SplitFunction extends TableFunction<Row> {
public void eval(String str, String delimiter) {
if (str == null) {
return;
}
for (String s : str.split(delimiter)) {
// 输出 Row
collect(Row.of(s, s.length()));
}
}
}
// 注册并使用
tableEnv.createTemporarySystemFunction("split", SplitFunction.class);
// LATERAL TABLE 语法
tableEnv.sqlQuery(
"SELECT user_id, word, length " +
"FROM users, LATERAL TABLE(split(tags, ',')) AS T(word, length)"
);
// 或使用 CROSS JOIN
tableEnv.sqlQuery(
"SELECT user_id, word, length " +
"FROM users " +
"CROSS JOIN LATERAL TABLE(split(tags, ',')) AS T(word, length)"
);
聚合函数 (UDAF)
聚合函数用于 GROUP BY 查询。
import org.apache.flink.table.functions.AggregateFunction;
// 定义累加器
public static class WeightedAvgAccum {
public double sum = 0;
public int count = 0;
}
// 加权平均函数
public class WeightedAvg extends AggregateFunction<Double, WeightedAvgAccum> {
@Override
public WeightedAvgAccum createAccumulator() {
return new WeightedAvgAccum();
}
@Override
public Double getValue(WeightedAvgAccum acc) {
if (acc.count == 0) {
return null;
}
return acc.sum / acc.count;
}
public void accumulate(WeightedAvgAccum acc, Double value, Integer weight) {
if (value != null && weight != null) {
acc.sum += value * weight;
acc.count += weight;
}
}
public void retract(WeightedAvgAccum acc, Double value, Integer weight) {
if (value != null && weight != null) {
acc.sum -= value * weight;
acc.count -= weight;
}
}
public void merge(WeightedAvgAccum acc, Iterable<WeightedAvgAccum> it) {
for (WeightedAvgAccum a : it) {
acc.sum += a.sum;
acc.count += a.count;
}
}
}
// 注册并使用
tableEnv.createTemporarySystemFunction("weightedAvg", WeightedAvg.class);
tableEnv.sqlQuery(
"SELECT user_id, weightedAvg(score, weight) as weighted_score " +
"FROM user_scores " +
"GROUP BY user_id"
);
连接器与格式
内置连接器
1. Kafka 连接器
-- Source
CREATE TABLE kafka_source (
`event_time` TIMESTAMP(3) METADATA FROM 'timestamp',
`partition` BIGINT METADATA VIRTUAL,
`offset` BIGINT METADATA VIRTUAL,
`user_id` STRING,
`behavior` STRING,
`item_id` STRING,
`ts` TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink-sql',
'scan.startup.mode' = 'latest-offset',
'format' = 'json'
);
-- Sink (支持 Exactly-Once)
CREATE TABLE kafka_sink (
user_id STRING,
cnt BIGINT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'user_stats',
'properties.bootstrap.servers' = 'localhost:9092',
'key.format' = 'json',
'value.format' = 'json'
);
2. JDBC 连接器
CREATE TABLE jdbc_table (
id INT,
name STRING,
age INT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/test',
'table-name' = 'users',
'driver' = 'com.mysql.cj.jdbc.Driver',
'username' = 'root',
'password' = 'password',
'sink.buffer-flush.max-rows' = '1000',
'sink.buffer-flush.interval' = '1s',
'sink.max-retries' = '3'
);
3. 文件系统连接器
CREATE TABLE fs_table (
user_id STRING,
order_amount DOUBLE,
dt STRING,
hr STRING
) PARTITIONED BY (dt, hr) WITH (
'connector' = 'filesystem',
'path' = 'hdfs://namenode:9000/data/orders',
'format' = 'parquet',
'sink.partition-commit.trigger' = 'process-time',
'sink.partition-commit.delay' = '1 hour',
'sink.partition-commit.policy.kind' = 'success-file'
);
自定义连接器
创建自定义 Source 连接器:
import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;
import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.RowData;
public class CustomTableSource extends RichParallelSourceFunction<RowData> {
private volatile boolean isRunning = true;
@Override
public void run(SourceContext<RowData> ctx) throws Exception {
while (isRunning) {
// 生成数据
GenericRowData rowData = new GenericRowData(3);
rowData.setField(0, "id_" + System.currentTimeMillis());
rowData.setField(1, "data_" + Math.random());
rowData.setField(2, System.currentTimeMillis());
ctx.collect(rowData);
Thread.sleep(1000);
}
}
@Override
public void cancel() {
isRunning = false;
}
}
// 创建 TableSource
public class CustomDynamicTableSource implements ScanTableSource {
@Override
public DynamicTableSource.DataStreamScanProvider getScanRuntimeProvider(
ScanTableSource.ScanContext scanContext) {
return new DataStreamScanProvider() {
@Override
public DataStream<RowData> produceDataStream(StreamExecutionEnvironment env) {
return env.addSource(new CustomTableSource());
}
};
}
@Override
public DynamicTableSource copy() {
return new CustomDynamicTableSource();
}
@Override
public String asSummaryString() {
return "Custom Table Source";
}
}
数据格式
内置格式
-- JSON 格式
CREATE TABLE json_table (
id INT,
data ROW<name STRING, age INT>
) WITH (
'connector' = 'kafka',
'topic' = 'test',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json',
'json.fail-on-missing-field' = 'false',
'json.ignore-parse-errors' = 'true'
);
-- Avro 格式
CREATE TABLE avro_table (
id INT,
data STRING
) WITH (
'connector' = 'kafka',
'topic' = 'test',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'avro',
'avro-confluent.url' = 'http://localhost:8081'
);
-- CSV 格式
CREATE TABLE csv_table (
id INT,
name STRING,
age INT
) WITH (
'connector' = 'filesystem',
'path' = '/path/to/csv',
'format' = 'csv',
'csv.field-delimiter' = ',',
'csv.ignore-parse-errors' = 'true'
);
实时数仓架构
分层设计
-- ODS 层 (原始数据层)
CREATE TABLE ods_user_behavior (
user_id STRING,
item_id STRING,
category_id STRING,
behavior STRING,
ts TIMESTAMP(3),
proctime AS PROCTIME(),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
-- DWD 层 (明细数据层)
CREATE TABLE dwd_user_behavior (
user_id STRING,
item_id STRING,
category_id STRING,
behavior STRING,
ts TIMESTAMP(3),
-- 清洗和标准化
behavior_type INT,
is_valid BOOLEAN
) WITH (
'connector' = 'kafka',
'topic' = 'dwd_user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
-- 数据清洗 ETL
INSERT INTO dwd_user_behavior
SELECT
user_id,
item_id,
category_id,
behavior,
ts,
CASE behavior
WHEN 'pv' THEN 1
WHEN 'buy' THEN 2
WHEN 'cart' THEN 3
WHEN 'fav' THEN 4
ELSE 0
END as behavior_type,
user_id IS NOT NULL AND item_id IS NOT NULL as is_valid
FROM ods_user_behavior
WHERE user_id IS NOT NULL;
-- DWS 层 (汇总数据层)
CREATE TABLE dws_user_stats_hourly (
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
user_id STRING,
pv_count BIGINT,
buy_count BIGINT,
cart_count BIGINT,
fav_count BIGINT,
PRIMARY KEY (window_start, user_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/dws',
'table-name' = 'user_stats_hourly',
'username' = 'root',
'password' = 'password'
);
-- 小时级汇总
INSERT INTO dws_user_stats_hourly
SELECT
TUMBLE_START(ts, INTERVAL '1' HOUR) as window_start,
TUMBLE_END(ts, INTERVAL '1' HOUR) as window_end,
user_id,
COUNT(CASE WHEN behavior = 'pv' THEN 1 END) as pv_count,
COUNT(CASE WHEN behavior = 'buy' THEN 1 END) as buy_count,
COUNT(CASE WHEN behavior = 'cart' THEN 1 END) as cart_count,
COUNT(CASE WHEN behavior = 'fav' THEN 1 END) as fav_count
FROM dwd_user_behavior
WHERE is_valid = true
GROUP BY
user_id,
TUMBLE(ts, INTERVAL '1' HOUR);
-- ADS 层 (应用数据层)
CREATE TABLE ads_user_activity_dashboard (
stat_date STRING,
active_users BIGINT,
total_pv BIGINT,
total_buy BIGINT,
conversion_rate DOUBLE,
PRIMARY KEY (stat_date) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://localhost:9200',
'index' = 'user_activity_dashboard'
);
维表关联
-- 用户维表
CREATE TABLE dim_users (
user_id STRING,
user_name STRING,
age INT,
gender STRING,
city STRING,
register_time TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/dim',
'table-name' = 'users',
'username' = 'root',
'password' = 'password',
'lookup.cache.max-rows' = '10000',
'lookup.cache.ttl' = '10 min'
);
-- 商品维表
CREATE TABLE dim_items (
item_id STRING,
item_name STRING,
category_id STRING,
price DOUBLE,
brand STRING,
PRIMARY KEY (item_id) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
'table-name' = 'dim_items',
'zookeeper.quorum' = 'localhost:2181'
);
-- 维表关联查询
CREATE VIEW user_behavior_enriched AS
SELECT
ub.*,
u.user_name,
u.age,
u.city,
i.item_name,
i.price,
i.brand
FROM dwd_user_behavior AS ub
JOIN dim_users FOR SYSTEM_TIME AS OF ub.proctime AS u
ON ub.user_id = u.user_id
JOIN dim_items FOR SYSTEM_TIME AS OF ub.proctime AS i
ON ub.item_id = i.item_id;
实时 ETL
public class RealtimeETLJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 配置 Checkpoint
env.enableCheckpointing(60000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 创建源表
tableEnv.executeSql(
"CREATE TABLE source_orders (" +
" order_id STRING," +
" user_id STRING," +
" amount DOUBLE," +
" currency STRING," +
" order_time TIMESTAMP(3)," +
" WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND" +
") WITH (" +
" 'connector' = 'kafka'," +
" 'topic' = 'orders'," +
" 'properties.bootstrap.servers' = 'localhost:9092'," +
" 'format' = 'json'" +
")"
);
// 汇率维表
tableEnv.executeSql(
"CREATE TABLE dim_currency_rates (" +
" currency STRING," +
" rate_to_usd DOUBLE," +
" update_time TIMESTAMP(3)," +
" PRIMARY KEY (currency) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'jdbc'," +
" 'url' = 'jdbc:mysql://localhost:3306/dim'," +
" 'table-name' = 'currency_rates'," +
" 'username' = 'root'," +
" 'password' = 'password'" +
")"
);
// 目标表
tableEnv.executeSql(
"CREATE TABLE sink_orders_normalized (" +
" order_id STRING," +
" user_id STRING," +
" amount_usd DOUBLE," +
" order_time TIMESTAMP(3)," +
" PRIMARY KEY (order_id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'upsert-kafka'," +
" 'topic' = 'orders_normalized'," +
" 'properties.bootstrap.servers' = 'localhost:9092'," +
" 'key.format' = 'json'," +
" 'value.format' = 'json'" +
")"
);
// ETL 逻辑
tableEnv.executeSql(
"INSERT INTO sink_orders_normalized " +
"SELECT " +
" o.order_id," +
" o.user_id," +
" o.amount * c.rate_to_usd as amount_usd," +
" o.order_time " +
"FROM source_orders AS o " +
"JOIN dim_currency_rates FOR SYSTEM_TIME AS OF o.order_time AS c " +
"ON o.currency = c.currency"
);
}
}
性能优化
1. MiniBatch 优化
Configuration configuration = tableEnv.getConfig().getConfiguration();
// 启用 MiniBatch
configuration.setString("table.exec.mini-batch.enabled", "true");
configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
configuration.setString("table.exec.mini-batch.size", "5000");
2. 本地-全局聚合
// 开启本地-全局聚合优化
configuration.setString("table.optimizer.agg-phase-strategy", "TWO_PHASE");
// 设置本地聚合的最小比例
configuration.setString("table.optimizer.distinct-agg.split.enabled", "true");
3. 异步 I/O 优化
-- 为维表查询启用异步 I/O
CREATE TABLE dim_table (
...
) WITH (
'connector' = 'jdbc',
'lookup.cache.max-rows' = '10000',
'lookup.cache.ttl' = '10 minutes',
'lookup.max-retries' = '3'
);
4. 状态后端优化
// 使用 RocksDB 状态后端
env.setStateBackend(new RocksDBStateBackend("hdfs://namenode:9000/flink/checkpoints"));
// 启用增量 Checkpoint
env.getCheckpointConfig().enableUnalignedCheckpoints();
5. 并行度调优
// 设置全局并行度
configuration.setString("parallelism.default", "8");
// 为特定算子设置并行度
Table result = tableEnv.sqlQuery("SELECT ... FROM ...");
DataStream<Row> stream = tableEnv.toDataStream(result);
stream.setParallelism(16);
最佳实践
1. Schema 设计
-- 使用合适的数据类型
CREATE TABLE optimized_table (
-- 使用 BIGINT 而不是 STRING 存储 ID
user_id BIGINT,
-- 使用 TIMESTAMP 而不是 STRING 存储时间
event_time TIMESTAMP(3),
-- 使用 DECIMAL 存储金额
amount DECIMAL(10, 2),
-- 合理设置字符串长度
description VARCHAR(255),
-- 使用 ROW 类型存储嵌套结构
user_info ROW<name STRING, age INT, tags ARRAY<STRING>>
) WITH (...);
2. 水印策略
-- 合理设置水印延迟
CREATE TABLE event_table (
event_time TIMESTAMP(3),
-- 根据实际数据延迟设置
WATERMARK FOR event_time AS event_time - INTERVAL '30' SECOND
) WITH (...);
-- 处理乱序数据
CREATE TABLE out_of_order_table (
event_time TIMESTAMP(3),
-- 允许 5 分钟的乱序
WATERMARK FOR event_time AS event_time - INTERVAL '5' MINUTE
) WITH (...);
3. 查询优化
-- 避免笛卡尔积
-- 不好的写法
SELECT * FROM table1, table2 WHERE table1.id > table2.id;
-- 好的写法
SELECT * FROM table1
JOIN table2 ON table1.key = table2.key
WHERE table1.id > table2.id;
-- 提前过滤
-- 不好的写法
SELECT * FROM (
SELECT * FROM large_table
) WHERE user_id = '123';
-- 好的写法
SELECT * FROM large_table WHERE user_id = '123';
-- 合理使用子查询
-- 使用 WITH 子句提高可读性
WITH user_stats AS (
SELECT user_id, COUNT(*) as cnt
FROM orders
GROUP BY user_id
)
SELECT * FROM user_stats WHERE cnt > 100;
4. 状态管理
// 设置状态 TTL
configuration.setString("table.exec.state.ttl", "24 h");
// 为不同的状态设置不同的 TTL
tableEnv.executeSql(
"CREATE TABLE stateful_table (" +
" user_id STRING," +
" last_login TIMESTAMP(3)," +
" PRIMARY KEY (user_id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'upsert-kafka'," +
" 'topic' = 'user_state'," +
" 'properties.bootstrap.servers' = 'localhost:9092'," +
" 'key.format' = 'json'," +
" 'value.format' = 'json'," +
" 'sink.buffer-flush.max-rows' = '0'," + // 立即刷新
" 'sink.buffer-flush.interval' = '0s'" +
")"
);
常见问题
Q1: 如何处理迟到数据?
-- 设置允许的最大延迟
CREATE TABLE late_data_table (
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '10' MINUTE
) WITH (...);
-- 使用侧输出流收集迟到数据
CREATE VIEW late_events AS
SELECT * FROM events
WHERE event_time < WATERMARK_TIME(event_time);
Q2: 如何优化维表 JOIN 性能?
-- 1. 使用缓存
CREATE TABLE dim_with_cache (
...
) WITH (
'lookup.cache.max-rows' = '10000',
'lookup.cache.ttl' = '10 min'
);
-- 2. 使用异步查找
CREATE TABLE dim_async (
...
) WITH (
'lookup.async' = 'true',
'lookup.async.timeout' = '180 s'
);
-- 3. 预加载维表到状态
-- 适用于小型维表
Q3: 如何处理数据倾斜?
-- 1. 使用随机前缀打散
SELECT
CONCAT(CAST(RAND() * 10 AS STRING), '_', user_id) as scatter_key,
user_id,
SUM(amount) as total
FROM orders
GROUP BY
CONCAT(CAST(RAND() * 10 AS STRING), '_', user_id),
user_id;
-- 2. 两阶段聚合
-- 第一阶段:局部聚合
CREATE VIEW local_agg AS
SELECT
user_id,
window_start,
SUM(amount) as partial_sum
FROM orders
GROUP BY
user_id,
TUMBLE(event_time, INTERVAL '1' MINUTE);
-- 第二阶段:全局聚合
SELECT
user_id,
SUM(partial_sum) as total_sum
FROM local_agg
GROUP BY user_id;
Q4: 如何实现 Exactly-Once?
// 1. 启用 Checkpoint
env.enableCheckpointing(60000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 2. 使用支持事务的 Sink
tableEnv.executeSql(
"CREATE TABLE exactly_once_sink (" +
" ..." +
") WITH (" +
" 'connector' = 'kafka'," +
" 'topic' = 'output'," +
" 'properties.bootstrap.servers' = 'localhost:9092'," +
" 'properties.transaction.timeout.ms' = '900000'," + // 15分钟
" 'sink.semantic' = 'exactly-once'," +
" 'sink.transactional-id-prefix' = 'flink-sink'" +
")"
);
总结
核心要点
| 组件 | 作用 | 使用场景 |
|---|---|---|
| Table API | 声明式流批统一处理 | 复杂数据处理逻辑 |
| SQL | 标准 SQL 查询 | 数据分析、ETL |
| UDF/UDTF/UDAF | 自定义函数扩展 | 特殊业务逻辑 |
| Connector | 外部系统集成 | 数据源和目标 |
| Catalog | 元数据管理 | 多表管理、Schema 演进 |
✅ 最佳实践总结
- 合理设计表结构,选择合适的数据类型
- 正确设置时间属性和水印,处理乱序数据
- 使用维表缓存提升 JOIN 性能
- 启用各种优化选项,如 MiniBatch、本地聚合
- 监控作业指标,及时发现和解决问题
🚫 避免:
- 不要在流式查询中使用无界的 JOIN
- 不要忽略水印设置导致的数据丢失
- 不要在高吞吐场景使用同步 I/O
适用场景
- 实时数仓: ODS → DWD → DWS → ADS 分层架构
- 实时 ETL: 数据清洗、转换、关联
- 实时分析: 实时报表、实时大屏
- 流批一体: 同时处理实时和历史数据
相关文章
Table & SQL 系列
- Flink SQL 内置函数 - 详细函数参考
- FlinkSQL 简明教程 - 快速入门指南
- Flink SQL 优化 - 查询优化技巧
DataStream API 系列
- Flink DataStream API 高级用法 - DataStream 高级特性
- DataStream 与 Table 转换 - API 互操作
实战案例
- Flink 实时数仓实战 - 完整案例
- Flink CDC - 实时数据同步
- Flink + Kafka 集成 - 端到端方案