Flink SQL 与常用数据库整合及类型映射
·
Flink SQL 与常用数据库整合及类型映射详细教程
📚 完整讲解 Flink SQL 与 MySQL、Doris、Kafka 的整合使用及统一 Java 数据类型映射
目录
- 1. 概述
- 2. Flink SQL 数据类型体系
- 3. MySQL CDC 整合
- 4. Apache Doris 整合
- 5. Kafka 整合
- 6. 统一类型映射方案
- 7. 完整实战案例
- 8. 类型转换工具类
- 9. 最佳实践
1. 概述
1.1 常见数据源对比
| 数据源 | 用途 | 典型场景 | Flink 连接器 |
|---|---|---|---|
| MySQL | OLTP 关系型数据库 | 业务数据存储、CDC 变更捕获 | flink-connector-mysql-cdc |
| Doris | OLAP 分析型数据库 | 实时数仓、多维分析 | 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 | 变长字符串 |
STRING | String | 无限长字符串 |
BOOLEAN | Boolean | 布尔值 |
BINARY(n) | byte[] | 定长二进制 |
VARBINARY(n) | byte[] | 变长二进制 |
BYTES | byte[] | 无限长二进制 |
DECIMAL(p, s) | java.math.BigDecimal | 精确数值 |
TINYINT | Byte | 1 字节整数 |
SMALLINT | Short | 2 字节整数 |
INT / INTEGER | Integer | 4 字节整数 |
BIGINT | Long | 8 字节整数 |
FLOAT | Float | 单精度浮点 |
DOUBLE | Double | 双精度浮点 |
DATE | java.time.LocalDate | 日期 |
TIME(p) | java.time.LocalTime | 时间 |
TIMESTAMP(p) | java.time.LocalDateTime | 时间戳(无时区) |
TIMESTAMP_LTZ(p) | java.time.Instant | 时间戳(带时区) |
INTERVAL | java.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 类型 | 说明 |
|---|---|---|---|
TINYINT | TINYINT | Byte | -128 到 127 |
SMALLINT | SMALLINT | Short | -32768 到 32767 |
MEDIUMINT | INT | Integer | -8388608 到 8388607 |
INT / INTEGER | INT | Integer | -2147483648 到 2147483647 |
BIGINT | BIGINT | Long | |
FLOAT | FLOAT | Float | 单精度浮点 |
DOUBLE | DOUBLE | Double | 双精度浮点 |
DECIMAL(p, s) | DECIMAL(p, s) | BigDecimal | 精确小数 |
BOOLEAN / TINYINT(1) | BOOLEAN | Boolean | |
CHAR(n) | CHAR(n) | String | 定长字符串 |
VARCHAR(n) | VARCHAR(n) | String | 变长字符串 |
TEXT | STRING | String | 长文本 |
BINARY(n) | BINARY(n) | byte[] | 定长二进制 |
VARBINARY(n) | VARBINARY(n) | byte[] | 变长二进制 |
BLOB | BYTES | byte[] | 二进制大对象 |
DATE | DATE | LocalDate | 日期 |
TIME(p) | TIME(p) | LocalTime | 时间 |
DATETIME(p) | TIMESTAMP(p) | LocalDateTime | 日期时间 |
TIMESTAMP(p) | TIMESTAMP_LTZ(p) | Instant | 带时区时间戳 |
YEAR | INT | Integer | 年份 |
ENUM | STRING | String | 枚举值 |
SET | ARRAY<STRING> | String[] | 集合 |
JSON | STRING | String | JSON 字符串 |
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 SQL 类型 | Java 类型 | 说明 |
|---|---|---|---|
BOOLEAN | BOOLEAN | Boolean | |
TINYINT | TINYINT | Byte | |
SMALLINT | SMALLINT | Short | |
INT | INT | Integer | |
BIGINT | BIGINT | Long | |
LARGEINT | STRING | String | 超大整数 |
FLOAT | FLOAT | Float | |
DOUBLE | DOUBLE | Double | |
DECIMAL(p, s) | DECIMAL(p, s) | BigDecimal | |
CHAR(n) | CHAR(n) | String | |
VARCHAR(n) | VARCHAR(n) | String | |
STRING | STRING | String | |
DATE | DATE | LocalDate | |
DATETIME | TIMESTAMP(0) | LocalDateTime | |
ARRAY<T> | ARRAY<T> | T[] | 数组 |
MAP<K,V> | MAP<K,V> | Map<K,V> | 映射 |
STRUCT<...> | ROW<...> | Row | 结构体 |
JSON | STRING | String | JSON 类型 |
JSONB | STRING | String | 二进制 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 类型 |
|---|---|---|
string | STRING | String |
number (整数) | BIGINT | Long |
number (小数) | DOUBLE | Double |
boolean | BOOLEAN | Boolean |
null | NULL | null |
array | ARRAY<T> | T[] |
object | ROW<...> | Row |
Avro 格式映射
| Avro 类型 | Flink SQL 类型 | Java 类型 |
|---|---|---|
boolean | BOOLEAN | Boolean |
int | INT | Integer |
long | BIGINT | Long |
float | FLOAT | Float |
double | DOUBLE | Double |
string | STRING | String |
bytes | BYTES | byte[] |
array | ARRAY<T> | T[] |
map | MAP<STRING, T> | Map<String, T> |
record | ROW<...> | Row |
enum | STRING | String |
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 |
|---|---|---|---|---|---|
| 订单ID | BIGINT | BIGINT | BIGINT | Long | number |
| 用户ID | BIGINT | BIGINT | BIGINT | Long | number |
| 商品名称 | VARCHAR(200) | VARCHAR(200) | STRING | String | string |
| 订单金额 | DECIMAL(10,2) | DECIMAL(10,2) | DECIMAL(10,2) | BigDecimal | number |
| 商品数量 | INT | INT | INT | Integer | number |
| 是否支付 | TINYINT(1) / BOOLEAN | BOOLEAN | BOOLEAN | Boolean | boolean |
| 订单状态 | TINYINT | TINYINT | TINYINT | Byte | number |
| 订单时间 | DATETIME | DATETIME | TIMESTAMP(0) | LocalDateTime | string |
| 创建时间 | TIMESTAMP | DATETIME | TIMESTAMP_LTZ(3) | Instant | string |
| 标签列表 | JSON | ARRAY<STRING> | ARRAY<STRING> | List<String> | array |
| 元数据 | JSON | MAP<STRING,STRING> | MAP<STRING,STRING> | Map<String,String> | object |
| 用户信息 | JSON | STRUCT<...> | ROW<...> | UserInfo | object |
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)
实现功能:
- 从 MySQL 实时捕获订单变更(CDC)
- 数据经过 Kafka 中转
- Flink SQL 进行清洗和聚合
- 写入 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 类型兼容性检查清单
| 检查项 | MySQL | Kafka JSON | Doris | Flink 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
✅ 实战案例:端到端数据管道实现
✅ 工具类封装:通用类型转换器
✅ 最佳实践:类型选择和性能优化建议
通过统一的类型映射方案,可以确保数据在不同系统间流转时的一致性和准确性!
更多推荐
所有评论(0)