全部笔记All notes

Flink Table & SQL API 实时数仓完整指南

阅读 9m 59s9m 59s read

概述

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 演进

✅ 最佳实践总结

  1. 合理设计表结构,选择合适的数据类型
  2. 正确设置时间属性和水印,处理乱序数据
  3. 使用维表缓存提升 JOIN 性能
  4. 启用各种优化选项,如 MiniBatch、本地聚合
  5. 监控作业指标,及时发现和解决问题

🚫 避免:

  • 不要在流式查询中使用无界的 JOIN
  • 不要忽略水印设置导致的数据丢失
  • 不要在高吞吐场景使用同步 I/O

适用场景

  • 实时数仓: ODS → DWD → DWS → ADS 分层架构
  • 实时 ETL: 数据清洗、转换、关联
  • 实时分析: 实时报表、实时大屏
  • 流批一体: 同时处理实时和历史数据

相关文章

Table & SQL 系列

DataStream API 系列

实战案例

上一章 / 下一章