Dan*_*Dan 5 java classloader apache-flink flink-streaming
我有一个源自Maven 入门项目的Flink 作业。该作业有一个打开 Postgres JDBC 连接的源。我正在使用示例docker-compose.yml在我自己的 Flink 会话集群上执行该作业。
当我第一次提交作业时,它执行成功。当我尝试再次提交时,出现以下错误:
Caused by: java.sql.SQLException: No suitable driver found for jdbc:postgresql://host.docker.internal:5432/postgres?user=postgres&password=mypassword
at java.sql.DriverManager.getConnection(DriverManager.java:689)
at java.sql.DriverManager.getConnection(DriverManager.java:270)
at com.myorg.project.JdbcPollingSource.run(JdbcPollingSource.java:25)
at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:110)
at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:66)
at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:269)
Run Code Online (Sandbox Code Playgroud)
我必须重新启动集群才能重新运行我的作业。为什么会发生这种情况?如何在不重启集群的情况下再次提交作业?
Maven 入门项目的唯一补充是:
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<version>42.2.24</version>
</dependency>
Run Code Online (Sandbox Code Playgroud)
Flink 源代码除了打开 JDBC 连接之外什么也不做,如下所示:
package com.mycompany;
import org.apache.flink.streaming.api.functions.source.RichSourceFunction;
import java.sql.Connection;
import java.sql.DriverManager;
public class JdbcSource extends RichSourceFunction<Integer> {
private final String connString;
public JdbcSource(String connString) {
this.connString = connString;
}
@Override
public void run(SourceContext<Integer> ctx) throws Exception {
try (Connection conn = DriverManager.getConnection(this.connString)) {
}
}
@Override
public void cancel() {
}
}
Run Code Online (Sandbox Code Playgroud)
我在 Flink 1.14.0 和 1.13.2 版本上进行了测试,结果相同。
请注意,这个问题提供了Class.forName("org.postgresql.Driver");在我的RichSourceFunction. 不过我想知道发生了什么事。
您可以参考的第一个问题是在 Apache Flink 中从 SQL 数据库读取 DataSet 时找不到 JDBC 驱动程序。
其次,如果使用会话模式。无需重新启动集群即可轻松重新运行 Flink 作业。您可以登录作业管理器 shell,然后使用命令重新运行作业。
Class.forName("org.postgresql.Driver");将触发静态方法块,因此您DriverManager可以获得驱动程序类。看:
// from org.postgresql.Driver
static {
try {
register();
} catch (SQLException var1) {
throw new ExceptionInInitializerError(var1);
}
}
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
1282 次 |
| 最近记录: |