Flink SQL 与常用数据库整合及类型映射详细教程

📚 完整讲解 Flink SQL 与 MySQL、Doris、Kafka 的整合使用及统一 Java 数据类型映射

目录


1. 概述

1.1 常见数据源对比

数据源用途典型场景Flink 连接器
MySQLOLTP 关系型数据库业务数据存储、CDC 变更捕获flink-connector-mysql-cdc
DorisOLAP 分析型数据库实时数仓、多维分析flink-connector-doris
Kafka消息队列流式数据传输、事件总线flink-connector-kafka

1.2 整合架构

MySQL (OLTP)  ─┐
               ├─→  Flink SQL (ETL)  ─→  Doris (OLAP)
Kafka (Stream)─┘                         
                                          ↓
                                     实时数据分析

2. Flink SQL 数据类型体系

2.1 Flink SQL 核心数据类型

根据 Flink 官方文档,Flink SQL 支持以下数据类型:

Flink SQL 类型Java 类型说明
CHAR(n)String定长字符串
VARCHAR(n)String变长字符串
STRINGString无限长字符串
BOOLEANBoolean布尔值
BINARY(n)byte[]定长二进制
VARBINARY(n)byte[]变长二进制
BYTESbyte[]无限长二进制
DECIMAL(p, s)java.math.BigDecimal精确数值
TINYINTByte1 字节整数
SMALLINTShort2 字节整数
INT / INTEGERInteger4 字节整数
BIGINTLong8 字节整数
FLOATFloat单精度浮点
DOUBLEDouble双精度浮点
DATEjava.time.LocalDate日期
TIME(p)java.time.LocalTime时间
TIMESTAMP(p)java.time.LocalDateTime时间戳(无时区)
TIMESTAMP_LTZ(p)java.time.Instant时间戳(带时区)
INTERVALjava.time.Duration / Period时间间隔
ARRAY<T>T[]数组
MAP<K, V>java.util.Map<K, V>映射
ROW<...>org.apache.flink.types.Row行/结构体

2.2 Table API 中的数据类型表示

import org.apache.flink.table.api.*;
import static org.apache.flink.table.api.DataTypes.*;

// 基本类型
DataType stringType = STRING();
DataType intType = INT();
DataType bigintType = BIGINT();
DataType doubleType = DOUBLE();
DataType booleanType = BOOLEAN();
DataType decimalType = DECIMAL(10, 2);

// 时间类型
DataType dateType = DATE();
DataType timeType = TIME();
DataType timestampType = TIMESTAMP(3);  // 精度 3 位(毫秒)
DataType timestampLtzType = TIMESTAMP_LTZ(3);

// 复杂类型
DataType arrayType = ARRAY(INT());
DataType mapType = MAP(STRING(), INT());
DataType rowType = ROW(
    FIELD("id", INT()),
    FIELD("name", STRING()),
    FIELD("amount", DECIMAL(10, 2))
);

3. MySQL CDC 整合

3.1 Maven 依赖

<dependencies>
    <!-- Flink MySQL CDC 连接器 -->
    <dependency>
        <groupId>com.ververica</groupId>
        <artifactId>flink-connector-mysql-cdc</artifactId>
        <version>2.4.2</version>
    </dependency>
    
    <!-- Flink Table API -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-table-api-java-bridge</artifactId>
        <version>1.20.0</version>
    </dependency>
</dependencies>

3.2 MySQL 到 Flink SQL 类型映射

根据 MySQL CDC 官方文档

MySQL 类型Flink SQL 类型Java 类型说明
TINYINTTINYINTByte-128 到 127
SMALLINTSMALLINTShort-32768 到 32767
MEDIUMINTINTInteger-8388608 到 8388607
INT / INTEGERINTInteger-2147483648 到 2147483647
BIGINTBIGINTLong
FLOATFLOATFloat单精度浮点
DOUBLEDOUBLEDouble双精度浮点
DECIMAL(p, s)DECIMAL(p, s)BigDecimal精确小数
BOOLEAN / TINYINT(1)BOOLEANBoolean
CHAR(n)CHAR(n)String定长字符串
VARCHAR(n)VARCHAR(n)String变长字符串
TEXTSTRINGString长文本
BINARY(n)BINARY(n)byte[]定长二进制
VARBINARY(n)VARBINARY(n)byte[]变长二进制
BLOBBYTESbyte[]二进制大对象
DATEDATELocalDate日期
TIME(p)TIME(p)LocalTime时间
DATETIME(p)TIMESTAMP(p)LocalDateTime日期时间
TIMESTAMP(p)TIMESTAMP_LTZ(p)Instant带时区时间戳
YEARINTInteger年份
ENUMSTRINGString枚举值
SETARRAY<STRING>String[]集合
JSONSTRINGStringJSON 字符串

3.3 创建 MySQL CDC 表

-- MySQL 数据库中的表结构
CREATE TABLE mysql_orders (
    order_id BIGINT PRIMARY KEY AUTO_INCREMENT,
    user_id BIGINT NOT NULL,
    product_name VARCHAR(200),
    amount DECIMAL(10, 2),
    status TINYINT,
    is_paid BOOLEAN,
    order_time DATETIME,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    metadata JSON
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- Flink SQL 中创建 MySQL CDC 源表
CREATE TABLE mysql_orders_cdc (
    order_id BIGINT,
    user_id BIGINT,
    product_name STRING,
    amount DECIMAL(10, 2),
    status TINYINT,
    is_paid BOOLEAN,
    order_time TIMESTAMP(0),
    created_at TIMESTAMP_LTZ(3),
    metadata STRING,
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'localhost',
    'port' = '3306',
    'username' = 'root',
    'password' = 'password',
    'database-name' = 'test_db',
    'table-name' = 'mysql_orders',
    'server-time-zone' = 'Asia/Shanghai',
    'scan.incremental.snapshot.enabled' = 'true'
);

3.4 Java 代码示例

package com.example.flink.mysql;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.Table;
import org.apache.flink.types.Row;

public class MySQLCDCExample {
    
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 创建 MySQL CDC 源表
        tableEnv.executeSql(
            "CREATE TABLE mysql_orders (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DECIMAL(10, 2)," +
            "  status TINYINT," +
            "  is_paid BOOLEAN," +
            "  order_time TIMESTAMP(0)," +
            "  created_at TIMESTAMP_LTZ(3)," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'test_db'," +
            "  'table-name' = 'mysql_orders'" +
            ")"
        );
        
        // 查询数据
        Table resultTable = tableEnv.sqlQuery(
            "SELECT " +
            "  order_id," +
            "  user_id," +
            "  product_name," +
            "  amount," +
            "  CASE WHEN is_paid THEN '已支付' ELSE '未支付' END AS payment_status," +
            "  order_time" +
            " FROM mysql_orders" +
            " WHERE status = 1"
        );
        
        // 转换为 DataStream 并打印
        tableEnv.toChangelogStream(resultTable).print();
        
        env.execute("MySQL CDC Example");
    }
}

4. Apache Doris 整合

4.1 Maven 依赖

<dependency>
    <groupId>org.apache.doris</groupId>
    <artifactId>flink-doris-connector-1.20</artifactId>
    <version>1.6.2</version>
</dependency>

4.2 Doris 到 Flink SQL 类型映射

根据 Doris Flink Connector 文档

Doris 类型Flink SQL 类型Java 类型说明
BOOLEANBOOLEANBoolean
TINYINTTINYINTByte
SMALLINTSMALLINTShort
INTINTInteger
BIGINTBIGINTLong
LARGEINTSTRINGString超大整数
FLOATFLOATFloat
DOUBLEDOUBLEDouble
DECIMAL(p, s)DECIMAL(p, s)BigDecimal
CHAR(n)CHAR(n)String
VARCHAR(n)VARCHAR(n)String
STRINGSTRINGString
DATEDATELocalDate
DATETIMETIMESTAMP(0)LocalDateTime
ARRAY<T>ARRAY<T>T[]数组
MAP<K,V>MAP<K,V>Map<K,V>映射
STRUCT<...>ROW<...>Row结构体
JSONSTRINGStringJSON 类型
JSONBSTRINGString二进制 JSON

4.3 创建 Doris 表

-- Doris 数据库中的表结构
CREATE TABLE doris_orders (
    order_id BIGINT NOT NULL,
    user_id BIGINT,
    product_name VARCHAR(200),
    amount DECIMAL(10, 2),
    order_count INT,
    avg_price DOUBLE,
    tags ARRAY<STRING>,
    created_at DATETIME
) UNIQUE KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 10;
-- Flink SQL 中创建 Doris Sink 表
CREATE TABLE doris_orders_sink (
    order_id BIGINT,
    user_id BIGINT,
    product_name STRING,
    amount DECIMAL(10, 2),
    order_count INT,
    avg_price DOUBLE,
    tags ARRAY<STRING>,
    created_at TIMESTAMP(0),
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = 'localhost:8030',
    'table.identifier' = 'test_db.doris_orders',
    'username' = 'root',
    'password' = '',
    'sink.label-prefix' = 'doris_label',
    'sink.properties.format' = 'json',
    'sink.properties.read_json_by_line' = 'true',
    'sink.enable-2pc' = 'false',
    'sink.buffer-size' = '10240',
    'sink.buffer-count' = '3',
    'sink.max-retries' = '3'
);

4.4 Java 代码示例

package com.example.flink.doris;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;

public class DorisExample {
    
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 创建 Doris Sink 表
        tableEnv.executeSql(
            "CREATE TABLE doris_sink (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DECIMAL(10, 2)," +
            "  order_count INT," +
            "  avg_price DOUBLE," +
            "  created_at TIMESTAMP(0)," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'doris'," +
            "  'fenodes' = 'localhost:8030'," +
            "  'table.identifier' = 'test_db.orders'," +
            "  'username' = 'root'," +
            "  'password' = ''," +
            "  'sink.properties.format' = 'json'," +
            "  'sink.buffer-size' = '10240'" +
            ")"
        );
        
        env.execute("Doris Sink Example");
    }
}

5. Kafka 整合

5.1 Maven 依赖

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kafka</artifactId>
    <version>1.20.0</version>
</dependency>

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-json</artifactId>
    <version>1.20.0</version>
</dependency>

5.2 Kafka 数据格式到 Flink SQL 类型映射

Kafka 本身不定义数据类型,需要通过格式(Format)解析:

JSON 格式映射
JSON 类型Flink SQL 类型Java 类型
stringSTRINGString
number (整数)BIGINTLong
number (小数)DOUBLEDouble
booleanBOOLEANBoolean
nullNULLnull
arrayARRAY<T>T[]
objectROW<...>Row
Avro 格式映射
Avro 类型Flink SQL 类型Java 类型
booleanBOOLEANBoolean
intINTInteger
longBIGINTLong
floatFLOATFloat
doubleDOUBLEDouble
stringSTRINGString
bytesBYTESbyte[]
arrayARRAY<T>T[]
mapMAP<STRING, T>Map<String, T>
recordROW<...>Row
enumSTRINGString

5.3 创建 Kafka 表

-- Kafka 消息格式(JSON)
{
    "order_id": 10001,
    "user_id": 1001,
    "product_name": "iPhone 15",
    "amount": 7999.99,
    "quantity": 1,
    "is_paid": true,
    "order_time": "2024-01-15 10:30:00",
    "tags": ["电子产品", "手机"],
    "user_info": {
        "name": "张三",
        "city": "北京"
    }
}
-- Flink SQL 创建 Kafka 源表
CREATE TABLE kafka_orders (
    order_id BIGINT,
    user_id BIGINT,
    product_name STRING,
    amount DOUBLE,
    quantity INT,
    is_paid BOOLEAN,
    order_time STRING,
    tags ARRAY<STRING>,
    user_info ROW<name STRING, city STRING>,
    -- 计算列:字符串转时间戳
    event_time AS TO_TIMESTAMP(order_time, 'yyyy-MM-dd HH:mm:ss'),
    -- 定义 Watermark
    WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'orders',
    'properties.bootstrap.servers' = 'localhost:9092',
    'properties.group.id' = 'flink-consumer',
    'scan.startup.mode' = 'earliest-offset',
    'format' = 'json',
    'json.fail-on-missing-field' = 'false',
    'json.ignore-parse-errors' = 'true'
);

5.4 Java 代码示例

package com.example.flink.kafka;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;

public class KafkaExample {
    
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 创建 Kafka 源表
        tableEnv.executeSql(
            "CREATE TABLE kafka_source (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DOUBLE," +
            "  quantity INT," +
            "  is_paid BOOLEAN," +
            "  order_time STRING," +
            "  event_time AS TO_TIMESTAMP(order_time, 'yyyy-MM-dd HH:mm:ss')," +
            "  WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND" +
            ") WITH (" +
            "  'connector' = 'kafka'," +
            "  'topic' = 'orders'," +
            "  'properties.bootstrap.servers' = 'localhost:9092'," +
            "  'properties.group.id' = 'flink-group'," +
            "  'scan.startup.mode' = 'latest-offset'," +
            "  'format' = 'json'" +
            ")"
        );
        
        // 查询并打印
        tableEnv.executeSql("SELECT * FROM kafka_source").print();
        
        env.execute("Kafka Source Example");
    }
}

6. 统一类型映射方案

6.1 通用数据模型定义

package com.example.flink.model;

import java.io.Serializable;
import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;

/**
 * 统一订单数据模型
 * 适用于 MySQL、Doris、Kafka 等多种数据源
 */
public class UnifiedOrder implements Serializable {
    
    // 基本类型
    private Long orderId;              // BIGINT
    private Long userId;               // BIGINT
    private String productName;        // STRING/VARCHAR
    private BigDecimal amount;         // DECIMAL(10, 2)
    private Integer quantity;          // INT
    private Boolean isPaid;            // BOOLEAN
    private Byte status;               // TINYINT
    
    // 时间类型
    private LocalDateTime orderTime;   // TIMESTAMP
    private LocalDateTime createdAt;   // TIMESTAMP_LTZ
    
    // 复杂类型
    private List<String> tags;         // ARRAY<STRING>
    private Map<String, String> metadata; // MAP<STRING, STRING>
    private UserInfo userInfo;         // ROW<...>
    
    // 嵌套类型
    public static class UserInfo implements Serializable {
        private String name;
        private String city;
        private Integer age;
        
        // Getters and Setters
        public String getName() { return name; }
        public void setName(String name) { this.name = name; }
        
        public String getCity() { return city; }
        public void setCity(String city) { this.city = city; }
        
        public Integer getAge() { return age; }
        public void setAge(Integer age) { this.age = age; }
    }
    
    // Getters and Setters
    public Long getOrderId() { return orderId; }
    public void setOrderId(Long orderId) { this.orderId = orderId; }
    
    public Long getUserId() { return userId; }
    public void setUserId(Long userId) { this.userId = userId; }
    
    public String getProductName() { return productName; }
    public void setProductName(String productName) { this.productName = productName; }
    
    public BigDecimal getAmount() { return amount; }
    public void setAmount(BigDecimal amount) { this.amount = amount; }
    
    public Integer getQuantity() { return quantity; }
    public void setQuantity(Integer quantity) { this.quantity = quantity; }
    
    public Boolean getIsPaid() { return isPaid; }
    public void setIsPaid(Boolean isPaid) { this.isPaid = isPaid; }
    
    public Byte getStatus() { return status; }
    public void setStatus(Byte status) { this.status = status; }
    
    public LocalDateTime getOrderTime() { return orderTime; }
    public void setOrderTime(LocalDateTime orderTime) { this.orderTime = orderTime; }
    
    public LocalDateTime getCreatedAt() { return createdAt; }
    public void setCreatedAt(LocalDateTime createdAt) { this.createdAt = createdAt; }
    
    public List<String> getTags() { return tags; }
    public void setTags(List<String> tags) { this.tags = tags; }
    
    public Map<String, String> getMetadata() { return metadata; }
    public void setMetadata(Map<String, String> metadata) { this.metadata = metadata; }
    
    public UserInfo getUserInfo() { return userInfo; }
    public void setUserInfo(UserInfo userInfo) { this.userInfo = userInfo; }
    
    @Override
    public String toString() {
        return "UnifiedOrder{" +
                "orderId=" + orderId +
                ", userId=" + userId +
                ", productName='" + productName + '\'' +
                ", amount=" + amount +
                ", quantity=" + quantity +
                ", isPaid=" + isPaid +
                ", status=" + status +
                ", orderTime=" + orderTime +
                ", createdAt=" + createdAt +
                ", tags=" + tags +
                ", metadata=" + metadata +
                ", userInfo=" + userInfo +
                '}';
    }
}

6.2 统一类型映射表

业务含义MySQL 类型Doris 类型Flink SQL 类型Java 类型Kafka JSON
订单IDBIGINTBIGINTBIGINTLongnumber
用户IDBIGINTBIGINTBIGINTLongnumber
商品名称VARCHAR(200)VARCHAR(200)STRINGStringstring
订单金额DECIMAL(10,2)DECIMAL(10,2)DECIMAL(10,2)BigDecimalnumber
商品数量INTINTINTIntegernumber
是否支付TINYINT(1) / BOOLEANBOOLEANBOOLEANBooleanboolean
订单状态TINYINTTINYINTTINYINTBytenumber
订单时间DATETIMEDATETIMETIMESTAMP(0)LocalDateTimestring
创建时间TIMESTAMPDATETIMETIMESTAMP_LTZ(3)Instantstring
标签列表JSONARRAY<STRING>ARRAY<STRING>List<String>array
元数据JSONMAP<STRING,STRING>MAP<STRING,STRING>Map<String,String>object
用户信息JSONSTRUCT<...>ROW<...>UserInfoobject

6.3 统一表结构定义

-- 统一的 Flink SQL 表定义(适用于多种数据源)
CREATE TABLE unified_orders (
    order_id BIGINT,                              -- 订单ID
    user_id BIGINT,                               -- 用户ID
    product_name STRING,                          -- 商品名称
    amount DECIMAL(10, 2),                        -- 订单金额
    quantity INT,                                 -- 商品数量
    is_paid BOOLEAN,                              -- 是否支付
    status TINYINT,                               -- 订单状态
    order_time TIMESTAMP(0),                      -- 订单时间(无时区)
    created_at TIMESTAMP_LTZ(3),                  -- 创建时间(带时区)
    tags ARRAY<STRING>,                           -- 标签数组
    metadata MAP<STRING, STRING>,                 -- 元数据映射
    user_info ROW<                                -- 用户信息(嵌套结构)
        name STRING,
        city STRING,
        age INT
    >,
    PRIMARY KEY (order_id) NOT ENFORCED
);

7. 完整实战案例

7.1 场景说明

业务流程:

MySQL (OLTP) → CDC → Kafka → Flink SQL ETL → Doris (OLAP)

实现功能:

  1. 从 MySQL 实时捕获订单变更(CDC)
  2. 数据经过 Kafka 中转
  3. Flink SQL 进行清洗和聚合
  4. 写入 Doris 供实时分析

7.2 完整代码实现

package com.example.flink.pipeline;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.Table;

/**
 * MySQL -> Kafka -> Flink SQL -> Doris 完整数据管道
 */
public class UnifiedDataPipeline {
    
    public static void main(String[] args) throws Exception {
        // 1. 创建执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);
        env.enableCheckpointing(60000); // 60秒 Checkpoint
        
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 2. 创建 MySQL CDC 源表
        createMySQLCDCTable(tableEnv);
        
        // 3. 创建 Kafka 中间表
        createKafkaTable(tableEnv);
        
        // 4. 创建 Doris Sink 表
        createDorisSinkTable(tableEnv);
        
        // 5. MySQL CDC -> Kafka
        tableEnv.executeSql(
            "INSERT INTO kafka_orders " +
            "SELECT " +
            "  order_id," +
            "  user_id," +
            "  product_name," +
            "  amount," +
            "  quantity," +
            "  is_paid," +
            "  status," +
            "  order_time," +
            "  created_at," +
            "  tags," +
            "  metadata," +
            "  user_info " +
            "FROM mysql_orders_source"
        );
        
        // 6. Kafka -> 实时聚合 -> Doris
        tableEnv.executeSql(
            "INSERT INTO doris_order_stats " +
            "SELECT " +
            "  user_id," +
            "  COUNT(*) as order_count," +
            "  SUM(amount) as total_amount," +
            "  AVG(amount) as avg_amount," +
            "  MAX(amount) as max_amount," +
            "  MIN(amount) as min_amount," +
            "  SUM(CASE WHEN is_paid THEN 1 ELSE 0 END) as paid_count," +
            "  CURRENT_TIMESTAMP as update_time " +
            "FROM kafka_orders " +
            "GROUP BY user_id"
        );
        
        env.execute("Unified Data Pipeline");
    }
    
    /**
     * 创建 MySQL CDC 源表
     */
    private static void createMySQLCDCTable(StreamTableEnvironment tableEnv) {
        tableEnv.executeSql(
            "CREATE TABLE mysql_orders_source (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DECIMAL(10, 2)," +
            "  quantity INT," +
            "  is_paid BOOLEAN," +
            "  status TINYINT," +
            "  order_time TIMESTAMP(0)," +
            "  created_at TIMESTAMP_LTZ(3)," +
            "  tags STRING," +  // MySQL JSON 存储为字符串
            "  metadata STRING," +
            "  user_info STRING," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'test_db'," +
            "  'table-name' = 'orders'," +
            "  'server-time-zone' = 'Asia/Shanghai'," +
            "  'scan.incremental.snapshot.enabled' = 'true'" +
            ")"
        );
    }
    
    /**
     * 创建 Kafka 中间表
     */
    private static void createKafkaTable(StreamTableEnvironment tableEnv) {
        tableEnv.executeSql(
            "CREATE TABLE kafka_orders (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DECIMAL(10, 2)," +
            "  quantity INT," +
            "  is_paid BOOLEAN," +
            "  status TINYINT," +
            "  order_time TIMESTAMP(0)," +
            "  created_at TIMESTAMP_LTZ(3)," +
            "  tags ARRAY<STRING>," +
            "  metadata MAP<STRING, STRING>," +
            "  user_info ROW<name STRING, city STRING, age INT>" +
            ") WITH (" +
            "  'connector' = 'kafka'," +
            "  'topic' = 'orders'," +
            "  'properties.bootstrap.servers' = 'localhost:9092'," +
            "  'properties.group.id' = 'flink-pipeline'," +
            "  'scan.startup.mode' = 'earliest-offset'," +
            "  'format' = 'json'," +
            "  'json.fail-on-missing-field' = 'false'" +
            ")"
        );
    }
    
    /**
     * 创建 Doris Sink 表
     */
    private static void createDorisSinkTable(StreamTableEnvironment tableEnv) {
        tableEnv.executeSql(
            "CREATE TABLE doris_order_stats (" +
            "  user_id BIGINT," +
            "  order_count BIGINT," +
            "  total_amount DECIMAL(20, 2)," +
            "  avg_amount DECIMAL(10, 2)," +
            "  max_amount DECIMAL(10, 2)," +
            "  min_amount DECIMAL(10, 2)," +
            "  paid_count BIGINT," +
            "  update_time TIMESTAMP(0)," +
            "  PRIMARY KEY (user_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'doris'," +
            "  'fenodes' = 'localhost:8030'," +
            "  'table.identifier' = 'dw.order_stats'," +
            "  'username' = 'root'," +
            "  'password' = ''," +
            "  'sink.properties.format' = 'json'," +
            "  'sink.buffer-size' = '10240'," +
            "  'sink.max-retries' = '3'" +
            ")"
        );
    }
}

7.3 端到端类型映射示例

package com.example.flink.pipeline;

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.Table;
import com.example.flink.model.UnifiedOrder;

/**
 * 端到端类型转换示例
 */
public class TypeConversionExample {
    
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 1. 从 MySQL CDC 读取数据
        tableEnv.executeSql(
            "CREATE TABLE mysql_source (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DECIMAL(10, 2)," +
            "  quantity INT," +
            "  is_paid BOOLEAN," +
            "  status TINYINT," +
            "  order_time TIMESTAMP(0)," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'test_db'," +
            "  'table-name' = 'orders'" +
            ")"
        );
        
        // 2. 将 Table 转换为 DataStream<UnifiedOrder>
        Table table = tableEnv.sqlQuery("SELECT * FROM mysql_source");
        DataStream<UnifiedOrder> orderStream = tableEnv.toDataStream(table, UnifiedOrder.class);
        
        // 3. 处理 DataStream
        DataStream<UnifiedOrder> processedStream = orderStream
            .filter(order -> order.getIsPaid())  // 只保留已支付订单
            .map(order -> {
                // 业务逻辑处理
                order.setStatus((byte) 2); // 更新状态
                return order;
            });
        
        // 4. 转换回 Table
        Table resultTable = tableEnv.fromDataStream(processedStream);
        tableEnv.createTemporaryView("processed_orders", resultTable);
        
        // 5. 写入 Doris
        tableEnv.executeSql(
            "CREATE TABLE doris_sink (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_name STRING," +
            "  amount DECIMAL(10, 2)," +
            "  quantity INT," +
            "  is_paid BOOLEAN," +
            "  status TINYINT," +
            "  order_time TIMESTAMP(0)," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'doris'," +
            "  'fenodes' = 'localhost:8030'," +
            "  'table.identifier' = 'dw.orders'," +
            "  'username' = 'root'," +
            "  'password' = ''" +
            ")"
        );
        
        tableEnv.executeSql("INSERT INTO doris_sink SELECT * FROM processed_orders");
        
        env.execute("Type Conversion Example");
    }
}

8. 类型转换工具类

8.1 通用类型转换器

package com.example.flink.util;

import org.apache.flink.table.data.*;
import org.apache.flink.table.types.logical.*;
import org.apache.flink.types.Row;

import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

/**
 * Flink SQL 类型转换工具类
 */
public class FlinkTypeConverter {
    
    /**
     * RowData 转 Java 对象
     */
    public static Object convertRowDataToJava(RowData rowData, int fieldIndex, LogicalType logicalType) {
        if (rowData.isNullAt(fieldIndex)) {
            return null;
        }
        
        LogicalTypeRoot typeRoot = logicalType.getTypeRoot();
        
        switch (typeRoot) {
            case CHAR:
            case VARCHAR:
                return rowData.getString(fieldIndex).toString();
            
            case BOOLEAN:
                return rowData.getBoolean(fieldIndex);
            
            case TINYINT:
                return rowData.getByte(fieldIndex);
            
            case SMALLINT:
                return rowData.getShort(fieldIndex);
            
            case INTEGER:
                return rowData.getInt(fieldIndex);
            
            case BIGINT:
                return rowData.getLong(fieldIndex);
            
            case FLOAT:
                return rowData.getFloat(fieldIndex);
            
            case DOUBLE:
                return rowData.getDouble(fieldIndex);
            
            case DECIMAL:
                DecimalType decimalType = (DecimalType) logicalType;
                return rowData.getDecimal(fieldIndex, decimalType.getPrecision(), decimalType.getScale())
                        .toBigDecimal();
            
            case DATE:
                int daysSinceEpoch = rowData.getInt(fieldIndex);
                return LocalDate.ofEpochDay(daysSinceEpoch);
            
            case TIME_WITHOUT_TIME_ZONE:
                int millisOfDay = rowData.getInt(fieldIndex);
                return LocalTime.ofNanoOfDay(millisOfDay * 1_000_000L);
            
            case TIMESTAMP_WITHOUT_TIME_ZONE:
                TimestampType timestampType = (TimestampType) logicalType;
                return rowData.getTimestamp(fieldIndex, timestampType.getPrecision())
                        .toLocalDateTime();
            
            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
                TimestampType timestampLtzType = (TimestampType) logicalType;
                return rowData.getTimestamp(fieldIndex, timestampLtzType.getPrecision())
                        .toInstant();
            
            case ARRAY:
                ArrayType arrayType = (ArrayType) logicalType;
                ArrayData arrayData = rowData.getArray(fieldIndex);
                return convertArrayDataToJava(arrayData, arrayType.getElementType());
            
            case MAP:
                MapType mapType = (MapType) logicalType;
                MapData mapData = rowData.getMap(fieldIndex);
                return convertMapDataToJava(mapData, mapType.getKeyType(), mapType.getValueType());
            
            case ROW:
                RowType rowType = (RowType) logicalType;
                RowData nestedRow = rowData.getRow(fieldIndex, rowType.getFieldCount());
                return convertRowDataToRow(nestedRow, rowType);
            
            case BINARY:
            case VARBINARY:
                return rowData.getBinary(fieldIndex);
            
            default:
                throw new UnsupportedOperationException("Unsupported type: " + typeRoot);
        }
    }
    
    /**
     * ArrayData 转 Java List
     */
    public static List<Object> convertArrayDataToJava(ArrayData arrayData, LogicalType elementType) {
        List<Object> list = new ArrayList<>();
        for (int i = 0; i < arrayData.size(); i++) {
            Object element = convertRowDataToJava(
                (RowData) arrayData,  // ArrayData 可以看作特殊的 RowData
                i,
                elementType
            );
            list.add(element);
        }
        return list;
    }
    
    /**
     * MapData 转 Java Map
     */
    public static Map<Object, Object> convertMapDataToJava(
            MapData mapData,
            LogicalType keyType,
            LogicalType valueType) {
        
        Map<Object, Object> map = new HashMap<>();
        ArrayData keyArray = mapData.keyArray();
        ArrayData valueArray = mapData.valueArray();
        
        for (int i = 0; i < mapData.size(); i++) {
            Object key = convertArrayElementToJava(keyArray, i, keyType);
            Object value = convertArrayElementToJava(valueArray, i, valueType);
            map.put(key, value);
        }
        
        return map;
    }
    
    /**
     * RowData 转 Row
     */
    public static Row convertRowDataToRow(RowData rowData, RowType rowType) {
        Row row = new Row(rowType.getFieldCount());
        for (int i = 0; i < rowType.getFieldCount(); i++) {
            Object value = convertRowDataToJava(rowData, i, rowType.getTypeAt(i));
            row.setField(i, value);
        }
        return row;
    }
    
    /**
     * 从 ArrayData 中获取元素
     */
    private static Object convertArrayElementToJava(ArrayData arrayData, int index, LogicalType type) {
        if (arrayData.isNullAt(index)) {
            return null;
        }
        
        switch (type.getTypeRoot()) {
            case VARCHAR:
            case CHAR:
                return arrayData.getString(index).toString();
            case BOOLEAN:
                return arrayData.getBoolean(index);
            case INTEGER:
                return arrayData.getInt(index);
            case BIGINT:
                return arrayData.getLong(index);
            case DOUBLE:
                return arrayData.getDouble(index);
            default:
                throw new UnsupportedOperationException("Unsupported array element type: " + type);
        }
    }
    
    /**
     * Java 对象转 RowData
     */
    public static GenericRowData convertJavaToRowData(Object[] values, LogicalType[] types) {
        GenericRowData rowData = new GenericRowData(values.length);
        
        for (int i = 0; i < values.length; i++) {
            if (values[i] == null) {
                rowData.setField(i, null);
                continue;
            }
            
            Object convertedValue = convertJavaValueToRowDataField(values[i], types[i]);
            rowData.setField(i, convertedValue);
        }
        
        return rowData;
    }
    
    /**
     * Java 值转 RowData 字段
     */
    private static Object convertJavaValueToRowDataField(Object value, LogicalType type) {
        if (value == null) {
            return null;
        }
        
        switch (type.getTypeRoot()) {
            case CHAR:
            case VARCHAR:
                return StringData.fromString(value.toString());
            
            case BOOLEAN:
                return value;
            
            case TINYINT:
            case SMALLINT:
            case INTEGER:
            case BIGINT:
            case FLOAT:
            case DOUBLE:
                return value;
            
            case DECIMAL:
                if (value instanceof BigDecimal) {
                    DecimalType decimalType = (DecimalType) type;
                    return DecimalData.fromBigDecimal(
                        (BigDecimal) value,
                        decimalType.getPrecision(),
                        decimalType.getScale()
                    );
                }
                break;
            
            case DATE:
                if (value instanceof LocalDate) {
                    return (int) ((LocalDate) value).toEpochDay();
                }
                break;
            
            case TIMESTAMP_WITHOUT_TIME_ZONE:
                if (value instanceof LocalDateTime) {
                    TimestampType timestampType = (TimestampType) type;
                    return TimestampData.fromLocalDateTime((LocalDateTime) value);
                }
                break;
            
            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
                if (value instanceof Instant) {
                    TimestampType timestampType = (TimestampType) type;
                    return TimestampData.fromInstant((Instant) value);
                }
                break;
            
            default:
                break;
        }
        
        return value;
    }
}

8.2 使用示例

package com.example.flink.demo;

import com.example.flink.util.FlinkTypeConverter;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.types.logical.*;

import java.math.BigDecimal;
import java.time.LocalDateTime;

public class TypeConverterDemo {
    
    public static void main(String[] args) {
        // 定义字段类型
        LogicalType[] types = new LogicalType[]{
            new BigIntType(),              // order_id
            new BigIntType(),              // user_id
            new VarCharType(200),          // product_name
            new DecimalType(10, 2),        // amount
            new IntType(),                 // quantity
            new BooleanType(),             // is_paid
            new TimestampType(0)           // order_time
        };
        
        // Java 对象数组
        Object[] values = new Object[]{
            10001L,
            1001L,
            "iPhone 15",
            new BigDecimal("7999.99"),
            1,
            true,
            LocalDateTime.now()
        };
        
        // Java -> RowData
        GenericRowData rowData = FlinkTypeConverter.convertJavaToRowData(values, types);
        System.out.println("RowData: " + rowData);
        
        // RowData -> Java
        for (int i = 0; i < types.length; i++) {
            Object javaValue = FlinkTypeConverter.convertRowDataToJava(rowData, i, types[i]);
            System.out.println("Field " + i + ": " + javaValue + " (" + javaValue.getClass().getSimpleName() + ")");
        }
    }
}

9. 最佳实践

9.1 类型选择建议

金额字段

-- 推荐使用 DECIMAL,避免精度丢失
amount DECIMAL(20, 2)  -- 而不是 DOUBLE

时间字段

-- 无时区场景
order_time TIMESTAMP(3)

-- 需要时区的场景
created_at TIMESTAMP_LTZ(3)

布尔字段

-- MySQL: TINYINT(1) 或 BOOLEAN
-- Flink SQL: BOOLEAN
is_active BOOLEAN

大整数

-- 超过 BIGINT 范围的数字
user_id_large STRING  -- 存储为字符串

9.2 类型兼容性检查清单

检查项MySQLKafka JSONDorisFlink SQL
整数范围✅ INT/BIGINT✅ number✅ INT/BIGINT✅ INT/BIGINT
小数精度✅ DECIMAL(p,s)⚠️ number✅ DECIMAL(p,s)✅ DECIMAL(p,s)
日期时间✅ DATETIME⚠️ string✅ DATETIME✅ TIMESTAMP
时区处理⚠️ 需配置❌ 不支持⚠️ 需配置✅ TIMESTAMP_LTZ
数组类型❌ 需JSON✅ array✅ ARRAY✅ ARRAY
嵌套结构❌ 需JSON✅ object✅ STRUCT✅ ROW

9.3 性能优化建议

合理选择精度

-- 不需要纳秒精度时,使用较低精度
TIMESTAMP(0)  -- 秒精度
TIMESTAMP(3)  -- 毫秒精度(推荐)
TIMESTAMP(6)  -- 微秒精度
TIMESTAMP(9)  -- 纳秒精度

避免不必要的类型转换

-- 不推荐
CAST(order_time AS STRING)

-- 推荐:源头统一类型
order_time TIMESTAMP(3)

使用合适的主键类型

-- BIGINT 比 STRING 更高效
PRIMARY KEY (order_id) NOT ENFORCED

10. 参考资源

官方文档链接


总结

本教程详细讲解了 Flink SQL 与 MySQL、Doris、Kafka 的整合及类型映射:

完整的类型映射表:涵盖所有常用数据类型
统一数据模型:跨系统的标准 Java POJO
实战案例:端到端数据管道实现
工具类封装:通用类型转换器
最佳实践:类型选择和性能优化建议

通过统一的类型映射方案,可以确保数据在不同系统间流转时的一致性和准确性!

更多推荐