如何让基于PostgreSQL的Apache Flink流作业在服务器持续运行
解决Flink JDBC作业无法持续运行的问题
嘿,这个问题我之前也踩过坑!你遇到的核心问题是:JDBCInputFormat是为批处理场景设计的输入格式,它只会一次性执行你的SELECT查询、拉取所有符合条件的数据,等这些数据处理完,作业就直接终止了,完全不会持续监听数据库里新增的数据。
要实现真正持续运行的流处理作业,有两种靠谱的方案,我给你详细拆解:
方案一:用Debezium CDC连接器(首推!)
这是目前处理数据库流数据的行业标准方案。Debezium能捕获PostgreSQL的WAL(预写日志),实时获取数据库的插入/更新/删除变更,真正实现流式的持续读取,完全符合你的需求。
操作步骤
- 先给PostgreSQL开启逻辑复制:修改
postgresql.conf配置,把wal_level设为logical,max_wal_senders设为至少10(具体数值根据你的集群规模调整),然后重启PostgreSQL服务。 - 在你的项目里引入Debezium PostgreSQL连接器的依赖(以Maven为例):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-debezium_${scala.binary.version}</artifactId> <version>${flink.version}</version> </dependency>
- 替换掉原来的
JDBCInputFormat,用CDC Source来读取数据,示例代码如下:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 配置Debezium连接参数,你可以把这些参数抽到配置类里 Properties props = new Properties(); props.setProperty("connector.class", "io.debezium.connector.postgresql.PostgresConnector"); props.setProperty("database.hostname", "你的PostgreSQL地址"); props.setProperty("database.port", "5432"); props.setProperty("database.user", "用户名"); props.setProperty("database.password", "密码"); props.setProperty("database.dbname", "你的数据库名"); props.setProperty("database.server.name", "postgres-server"); props.setProperty("table.include.list", "你的数据库名.源表名"); // 指定要监听的表 // 创建CDC Source DebeziumSource<String> cdcSource = DebeziumSource.<String>builder() .connectorProperties(props) .deserializer(new StringDebeziumDeserializationSchema()) // 这里可以自定义反序列化逻辑,把CDC事件转成你需要的Row类型 .build(); // 读取CDC流数据 DataStream<String> cdcStream = env.fromSource(cdcSource, WatermarkStrategy.noWatermarks(), "PostgreSQL CDC Source"); // 把CDC事件转换成Row类型,这里需要你根据实际的CDC事件格式来写map逻辑 SingleOutputStreamOperator<Row> sourceStream = cdcStream .map(jsonStr -> { // 示例:用FastJSON解析CDC的JSON字符串,转换成Row JSONObject json = JSON.parseObject(jsonStr); JSONObject after = json.getJSONObject("after"); Row row = new Row(3); row.setField(0, after.getLong("id")); row.setField(1, after.getString("name")); row.setField(2, after.getDate("create_time")); return row; }) .assignTimestampsAndWatermarks(WatermarkStrategy .<Row>forMonotonousTimestamps() .withTimestampAssigner((row, ts) -> ((Date) row.getField(2)).getTime())); // 后面的处理逻辑和你原来的代码完全兼容 SingleOutputStreamOperator<Row> processedStream = sourceStream .keyBy(0) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .fold(null, new FoldFunction<Row, Row>() { @Override public Row fold(Row row1, Row row) throws Exception { return row; } }); // 写入目标数据库,继续用你原来的JDBCOutputFormat就行 processedStream.writeUsingOutputFormat(JDBCOutputFormat.buildJDBCOutputFormat() .setDrivername(JDBCConfig.DRIVER_CLASS) .setDBUrl(JDBCConfig.DB_URL) .setQuery("insert into tablename(id, name) values (?,?)") .setSqlTypes(new int[]{Types.BIGINT, Types.VARCHAR}) .finish()); env.execute("PostgreSQL CDC Stream Job");
方案二:自定义周期性查询的Source(临时过渡方案)
如果暂时没法用CDC(比如数据库权限不够开启逻辑复制),可以自定义一个SourceFunction,每隔一段时间执行一次SELECT查询,只拉取上次查询之后新增的数据(前提是你的源表有自增ID或者时间戳字段来做过滤)。
自定义Source代码
public class PeriodicJDBCSource implements SourceFunction<Row> { private volatile boolean running = true; private final long queryIntervalMs; // 每次查询的间隔,比如5000ms private final String dbUrl; private final String driverClass; private final String queryTemplate; // 带过滤条件的查询语句,比如"SELECT * FROM source_table WHERE create_time > ?" private long lastQueryTimestamp = 0; // 记录上次查询的最大时间戳 public PeriodicJDBCSource(long queryIntervalMs, String dbUrl, String driverClass, String queryTemplate) { this.queryIntervalMs = queryIntervalMs; this.dbUrl = dbUrl; this.driverClass = driverClass; this.queryTemplate = queryTemplate; } @Override public void run(SourceContext<Row> ctx) throws Exception { Class.forName(driverClass); while (running) { try (Connection conn = DriverManager.getConnection(dbUrl); PreparedStatement stmt = conn.prepareStatement(queryTemplate)) { // 设置过滤条件:只查上次查询之后新增的数据 stmt.setTimestamp(1, new Timestamp(lastQueryTimestamp)); ResultSet rs = stmt.executeQuery(); // 把ResultSet转换成Row并发送到流中 while (rs.next()) { Row row = new Row(3); // 根据你的表结构调整字段数量 row.setField(0, rs.getLong("id")); row.setField(1, rs.getString("name")); row.setField(2, rs.getTimestamp("create_time")); ctx.collect(row); // 更新lastQueryTimestamp为当前数据的时间戳,避免重复读取 lastQueryTimestamp = Math.max(lastQueryTimestamp, rs.getTimestamp("create_time").getTime()); } } catch (SQLException e) { e.printStackTrace(); // 出现异常可以加个重试逻辑,避免直接终止 Thread.sleep(1000); } // 等待下一次查询 Thread.sleep(queryIntervalMs); } } @Override public void cancel() { running = false; } }
然后在主作业里使用这个自定义Source:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); // 使用自定义的周期性JDBC Source,每隔5秒查询一次 DataStream<Row> sourceStream = env.addSource(new PeriodicJDBCSource( 5000, JDBCConfig.DB_URL, JDBCConfig.DRIVER_CLASS, "SELECT * FROM source_table WHERE create_time > ?" )); // 后面的处理和写入逻辑和你原来的代码一致 SingleOutputStreamOperator<Row> processedStream = sourceStream .assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Row>() { @Override public long extractAscendingTimestamp(Row row) { Date dt = (Date) row.getField(2); return dt.getTime(); } }) .keyBy(0) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .fold(null, new FoldFunction<Row, Row>() { @Override public Row fold(Row row1, Row row) throws Exception { return row; } }); processedStream.writeUsingOutputFormat(JDBCOutputFormat.buildJDBCOutputFormat() .setDrivername(JDBCConfig.DRIVER_CLASS) .setDBUrl(JDBCConfig.DB_URL) .setQuery("insert into tablename(id, name) values (?,?)") .setSqlTypes(new int[]{Types.BIGINT, Types.VARCHAR}) .finish()); env.execute("Periodic JDBC Stream Job");
一些额外提醒
- 方案二的局限性很大:数据延迟至少是你设置的查询间隔,而且如果有数据更新频繁的情况,可能会出现重复读取或者漏读,只适合数据更新频率低的临时场景。
- 如果你用的是Flink 1.12及以上版本,建议用新的
WatermarkStrategy替代旧的AscendingTimestampExtractor,代码更简洁,也更符合Flink的新API规范。
内容的提问来源于stack exchange,提问作者Mul
相关产品推荐
相关产品推荐

