You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让基于PostgreSQL的Apache Flink流作业在服务器持续运行

嘿,这个问题我之前也踩过坑!你遇到的核心问题是:JDBCInputFormat是为批处理场景设计的输入格式,它只会一次性执行你的SELECT查询、拉取所有符合条件的数据,等这些数据处理完,作业就直接终止了,完全不会持续监听数据库里新增的数据。

要实现真正持续运行的流处理作业,有两种靠谱的方案,我给你详细拆解:


方案一:用Debezium CDC连接器(首推!)

这是目前处理数据库流数据的行业标准方案。Debezium能捕获PostgreSQL的WAL(预写日志),实时获取数据库的插入/更新/删除变更,真正实现流式的持续读取,完全符合你的需求。

操作步骤

  1. 先给PostgreSQL开启逻辑复制:修改postgresql.conf配置,把wal_level设为logical,max_wal_senders设为至少10(具体数值根据你的集群规模调整),然后重启PostgreSQL服务。
  2. 在你的项目里引入Debezium PostgreSQL连接器的依赖(以Maven为例):
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-debezium_${scala.binary.version}</artifactId>
    <version>${flink.version}</version>
</dependency>
  1. 替换掉原来的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 07:52:33