报错

2025-08-05 18:14:44
org.apache.flink.util.FlinkException: Global failure triggered by OperatorCoordinator for 'Source: PostgresParallelSource -> amos.time_captured-amos.billed_items-amos.pickslip_detail-amos.bill_details-amos.wp_content-amos.consumables-amos.cm_project-amos.staff_pqs_type-amos.ppl_plan-amos.phase-amos.rm_resource_requirement-amos.wp_sequence-amos.general_ledger-amos.ext_alloc_status_detail-amos.part_location-amos.evt_commercial_requirement-amos.cm_project_partition-amos.sh_header-amos.address-amos.contract_mat_scope_attrib-amos.contract_mat_scope-amos.ac_typ-amos.measure_unit-amos.project_type-amos.sh_detail: Writer -> amos.time_captured-amos.billed_items-amos.pickslip_detail-amos.bill_details-amos.wp_content-amos.consumables-amos.cm_project-amos.staff_pqs_type-amos.ppl_plan-amos.phase-amos.rm_resource_requirement-amos.wp_sequence-amos.general_ledger-amos.ext_alloc_status_detail-amos.part_location-amos.evt_commercial_requirement-amos.cm_project_partition-amos.sh_header-amos.address-amos.contract_mat_scope_attrib-amos.contract_mat_scope-amos.ac_typ-amos.measure_unit-amos.project_type-amos.sh_detail: Committer' (operator cbc357ccb763df2852fee8c4fc7d55f2).
    at org.apache.flink.runtime.operators.coordination.OperatorCoordinatorHolder$LazyInitializedCoordinatorContext.failJob(OperatorCoordinatorHolder.java:624)
    at org.apache.flink.runtime.operators.coordination.RecreateOnResetOperatorCoordinator$QuiesceableContext.failJob(RecreateOnResetOperatorCoordinator.java:248)
    at org.apache.flink.runtime.source.coordinator.SourceCoordinatorContext.failJob(SourceCoordinatorContext.java:395)
    at org.apache.flink.runtime.source.coordinator.SourceCoordinator.lambda$runInEventLoop$10(SourceCoordinator.java:483)
    at org.apache.flink.util.ThrowableCatchingRunnable.run(ThrowableCatchingRunnable.java:40)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source)
    at java.base/java.util.concurrent.FutureTask.run(Unknown Source)
    at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(Unknown Source)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
    at java.base/java.lang.Thread.run(Unknown Source)
Caused by: org.apache.flink.util.FlinkRuntimeException: Generate Splits for table amos.sh_header error
    at com.ververica.cdc.connectors.postgres.source.PostgresChunkSplitter.generateSplits(PostgresChunkSplitter.java:124)
    at com.ververica.cdc.connectors.base.source.assigner.SnapshotSplitAssigner.getNext(SnapshotSplitAssigner.java:178)
    at com.ververica.cdc.connectors.base.source.assigner.HybridSplitAssigner.getNext(HybridSplitAssigner.java:137)
    at com.ververica.cdc.connectors.base.source.enumerator.IncrementalSourceEnumerator.assignSplits(IncrementalSourceEnumerator.java:174)
    at com.ververica.cdc.connectors.base.source.enumerator.IncrementalSourceEnumerator.handleSplitRequest(IncrementalSourceEnumerator.java:97)
    at org.apache.flink.runtime.source.coordinator.SourceCoordinator.handleRequestSplitEvent(SourceCoordinator.java:568)
    at org.apache.flink.runtime.source.coordinator.SourceCoordinator.lambda$handleEventFromOperator$3(SourceCoordinator.java:295)
    at org.apache.flink.runtime.source.coordinator.SourceCoordinator.lambda$runInEventLoop$10(SourceCoordinator.java:469)
    ... 7 more
Caused by: org.apache.kafka.connect.errors.ConnectException: Database processing error
    at io.debezium.connector.postgresql.PostgresOffsetContext.initialContext(PostgresOffsetContext.java:251)
    at io.debezium.connector.postgresql.PostgresOffsetContext.initialContext(PostgresOffsetContext.java:228)
    at com.ververica.cdc.connectors.postgres.source.utils.CustomPostgresSchema.readTableSchema(CustomPostgresSchema.java:69)
    at com.ververica.cdc.connectors.postgres.source.utils.CustomPostgresSchema.getTableSchema(CustomPostgresSchema.java:58)
    at com.ververica.cdc.connectors.postgres.source.PostgresDialect.queryTableSchema(PostgresDialect.java:197)
    at com.ververica.cdc.connectors.postgres.source.PostgresChunkSplitter.generateSplits(PostgresChunkSplitter.java:90)
    ... 14 more
Caused by: org.postgresql.util.PSQLException: This connection has been closed.
    at org.postgresql.jdbc.PgConnection.checkClosed(PgConnection.java:1009)
    at org.postgresql.jdbc.PgConnection.getMetaData(PgConnection.java:1425)
    at io.debezium.connector.postgresql.connection.PostgresConnection.currentXLogLocation(PostgresConnection.java:554)
    at io.debezium.connector.postgresql.PostgresOffsetContext.initialContext(PostgresOffsetContext.java:235)
    ... 19 more

在这里插入图片描述
在这里插入图片描述
最后跳到这里

    public synchronized Connection connection(boolean executeOnConnect) throws SQLException {
        if (!isConnected()) {
            conn = factory.connect(JdbcConfiguration.adapt(config));
            if (!isConnected()) {
                throw new SQLException("Unable to obtain a JDBC connection");
            }
            // Always run the initial operations on this new connection
            if (initialOps != null) {
                execute(initialOps);
            }
            final String statements = config.getString(JdbcConfiguration.ON_CONNECT_STATEMENTS);
            if (statements != null && executeOnConnect) {
                final List<String> splitStatements = parseSqlStatementString(statements);
                execute(splitStatements.toArray(new String[splitStatements.size()]));
            }
        }
        return conn;
    }

if (!isConnected())客户端永远都是连接的,再去执行下面代码的时候其实是一个无效的链接,被关闭了,就会报错,所以这里应该用 if ( !isConnected() && isValid() ),并且要判断是否有效,没有效的话再
conn = factory.connect(JdbcConfiguration.adapt(config));
拿新的connect。

我还没有验证,仅仅是一个思路,后面补上 thx

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐