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. 关键配置项说明
参数推荐值作用
fetchSize500-5000控制单次查询数据量
maximumPoolSizeCPU核心数×2连接池最大连接数
idleTimeout30000ms空闲连接回收时间
connectionTimeout10000ms连接超时时间
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. 性能优化策略
  1. 索引优化:确保查询字段(如id)有索引
  2. 分片策略
    • 范围分片:WHERE id BETWEEN ? AND ?
    • 哈希分片:WHERE HASH(id) % N = M
  3. 增量查询:通过last_update_time字段增量拉取
  4. 批量提交:每累积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实时捕获变更,避免全表扫描

更多推荐