Flink 自定义 Source 函数开发:从 MySQL 读取数据的实操
·
Flink 自定义 Source 函数:MySQL 数据读取实操指南
1. 核心实现原理
- 并行读取:通过
RichParallelSourceFunction实现多线程并发读取 - 断点续传:结合
CheckpointedFunction保存offset状态 - 连接池管理:使用
HikariCP高效管理数据库连接 - 数据分片:根据主键范围或时间窗口分割查询任务
2. 依赖配置(Maven)
<dependencies>
<!-- Flink 核心依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.15.0</version>
</dependency>
<!-- MySQL 连接器 -->
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>8.0.30</version>
</dependency>
<!-- 连接池 -->
<dependency>
<groupId>com.zaxxer</groupId>
<artifactId>HikariCP</artifactId>
<version>5.0.1</version>
</dependency>
</dependencies>
3. 完整源码实现
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;
import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction;
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.runtime.state.FunctionInitializationContext;
import org.apache.flink.runtime.state.FunctionSnapshotContext;
import com.zaxxer.hikari.HikariDataSource;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
public class MySQLSource extends RichParallelSourceFunction<User>
implements CheckpointedFunction {
private volatile boolean isRunning = true;
private transient HikariDataSource dataSource;
private transient ListState<Long> offsetState;
private long currentOffset = 0;
// 配置参数
private final String jdbcUrl;
private final String username;
private final String password;
private final String tableName;
public MySQLSource(String jdbcUrl, String username,
String password, String tableName) {
this.jdbcUrl = jdbcUrl;
this.username = username;
this.password = password;
this.tableName = tableName;
}
@Override
public void open(Configuration parameters) {
// 初始化连接池
dataSource = new HikariDataSource();
dataSource.setJdbcUrl(jdbcUrl);
dataSource.setUsername(username);
dataSource.setPassword(password);
dataSource.setMaximumPoolSize(5);
}
@Override
public void run(SourceContext<User> ctx) throws Exception {
final int taskIdx = getRuntimeContext().getIndexOfThisSubtask();
final int totalTasks = getRuntimeContext().getNumberOfParallelSubtasks();
try (Connection conn = dataSource.getConnection()) {
String sql = String.format(
"SELECT id, name, create_time FROM %s " +
"WHERE MOD(id, %d) = %d AND id > ? " + // 分片策略
"ORDER BY id ASC", tableName, totalTasks, taskIdx);
PreparedStatement ps = conn.prepareStatement(sql);
ps.setFetchSize(1000); // 批量获取
while (isRunning) {
ps.setLong(1, currentOffset);
ResultSet rs = ps.executeQuery();
int count = 0;
while (rs.next() && isRunning) {
User user = new User(
rs.getLong("id"),
rs.getString("name"),
rs.getTimestamp("create_time").toLocalDateTime()
);
ctx.collect(user);
currentOffset = user.getId();
count++;
}
if (count == 0) {
Thread.sleep(5000); // 无数据时暂停
}
}
}
}
@Override
public void cancel() {
isRunning = false;
}
@Override
public void close() {
if (dataSource != null) {
dataSource.close();
}
}
// === 状态管理 ===
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
offsetState.clear();
offsetState.add(currentOffset);
}
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
ListStateDescriptor<Long> descriptor =
new ListStateDescriptor<>("offsetState", Long.class);
offsetState = context.getOperatorStateStore().getListState(descriptor);
if (context.isRestored()) {
for (Long offset : offsetState.get()) {
currentOffset = offset;
}
}
}
// 数据实体类
public static class User {
private long id;
private String name;
private LocalDateTime createTime;
// 构造方法/getters/setters省略
}
}
4. 关键配置项说明
| 参数 | 推荐值 | 作用 |
|---|---|---|
fetchSize | 500-5000 | 控制单次查询数据量 |
maximumPoolSize | CPU核心数×2 | 连接池最大连接数 |
idleTimeout | 30000ms | 空闲连接回收时间 |
connectionTimeout | 10000ms | 连接超时时间 |
5. Flink作业调用示例
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<User> userStream = env.addSource(
new MySQLSource(
"jdbc:mysql://localhost:3306/mydb",
"root",
"password123",
"user_table"
)
).setParallelism(4); // 设置并行度
userStream.print();
env.execute("MySQL Source Job");
6. 性能优化策略
- 索引优化:确保查询字段(如
id)有索引 - 分片策略:
- 范围分片:
WHERE id BETWEEN ? AND ? - 哈希分片:
WHERE HASH(id) % N = M
- 范围分片:
- 增量查询:通过
last_update_time字段增量拉取 - 批量提交:每累积1000条数据手动触发
checkpoint
7. 异常处理建议
try {
// 数据库操作
} catch (SQLException e) {
if (e.getErrorCode() == 1060) { // 连接超时
dataSource.getConnection().close(); // 重置连接
} else if (e.getErrorCode() == 2013) { // 查询中断
throw new RuntimeException("Query interrupted", e);
}
}
最佳实践:生产环境建议配合
Debezium实现CDC实时捕获变更,避免全表扫描
更多推荐
所有评论(0)