Tengo un trabajo de Flink derivado del proyecto inicial Maven . Ese trabajo tiene una fuente que abre una conexión JDBC de Postgres. Estoy ejecutando el trabajo en mi propio clúster de sesión de Flink usando el ejemplo docker-compose.yml .
Cuando envío el trabajo por primera vez, se ejecuta correctamente. Cuando intento enviarlo de nuevo me sale el siguiente error:
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)Tengo que reiniciar mi clúster para volver a ejecutar mi trabajo. ¿Por qué está pasando esto? ¿Cómo puedo volver a enviar mi trabajo sin tener que reiniciar el clúster?
La única adición al proyecto inicial de Maven es:
<dependency> <groupId>org.postgresql</groupId> <artifactId>postgresql</artifactId> <version>42.2.24</version> </dependency>La fuente de Flink no hace más que abrir una conexión JDBC y es la siguiente:
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() { } }He probado esto en Flink versión 1.14.0 y 1.13.2 con los mismos resultados.
Tenga en cuenta que esta pregunta proporciona una solución para usar Class.forName("org.postgresql.Driver"); dentro de mi RichSourceFunction . Sin embargo, me gustaría saber qué está pasando.
La primera pregunta a la que puede hacer referencia al controlador JDBC no se puede encontrar al leer un conjunto de datos de una base de datos SQL en Apache Flink .
En segundo lugar, si usa el modo de sesión. Puede ser fácil volver a ejecutar el trabajo de Flink sin reiniciar el clúster. puede iniciar sesión en el shell del administrador de trabajos y luego usar el comando volver a ejecutar el trabajo.
Class.forName("org.postgresql.Driver"); activará el bloque de método estático, por lo que DriverManager puede obtener la clase de controlador. ver:
// from org.postgresql.Driver static { try { register(); } catch (SQLException var1) { throw new ExceptionInInitializerError(var1); } }Tengo esta dependencia pom.xml para Postgres para Apache Flink 1.13:
<dependency> <groupId>org.postgresql</groupId> <artifactId>postgresql</artifactId> <version>9.4-1201-jdbc41</version> </dependency>puede tener una clase de conector Postgres, por ejemplo:
public class PostgreSQLConnector { private static volatile PostgreSQLConnector instance; private Connection connectionDB = null; public PostgreSQLConnector(your params) { ... } public static PostgreSQLConnector getInstance() { PostgreSQLConnector postgreSQLConnector = instance; if (postgreSQLConnector != null) return postgreSQLConnector; synchronized (PostgreSQLConnector.class) { if (instance == null) { instance = new PostgreSQLConnector(your params); } return instance; } } public Connection getConnectionDB() throws SQLException { if (checkNullConnection()) CreateConnection(); return connectionDB; } public void CheckConnection() throws SQLException { if (checkNullConnection()) CreateConnection(); } public void CreateConnection() throws SQLException { try { Class.forName(sink.driverName); connectionDB = DriverManager.getConnection(fullUrl, username, password); } catch (Exception e) { ... } } public boolean checkNullConnection() throws SQLException { return (connectionDB == null || connectionDB.isClosed()); } } luego puede crear una RichSourceFunction y crear la conexión en el método de open de anulaciones, no en la run
public class JdbcSource extends RichSourceFunction<Integer> { private final String connString; private static Connection dbConnection; private static final PostgreSQLConnector postgreSQLConnector = PostgreSQLConnector.getInstance(); public JdbcSource(String connString) { this.connString = connString; } @Override public void open(Configuration parameters) throws SQLException { dbConnection = postgreSQLConnector.getConnectionDB(); } @Override public void close() throws Exception { if (dbConnection != null) dbConnection.close(); } @Override public void run(SourceContext<Integer> ctx) throws Exception { do something here with the connection } @Override public void cancel() { } }Algo así podrías probar y debería funcionar
De acuerdo con la documentación oficial del controlador JDBC de PostgreSQL, si está utilizando Java 1.6+, puede colocar el archivo jar del controlador en el classpath. La JVM cargará el controlador automáticamente. Entonces, la pregunta es cómo colocar el archivo jar del controlador en el classpath.
Dado que está utilizando la ventana acoplable para implementar un clúster de sesión, hay dos formas en que puede funcionar:
Ejecute y acceda a la imagen con el comando:
docker docker run -it -v $PWD:/tmp/flink <address to image> -- bash Copie el archivo jar del controlador en la carpeta /opt/flink/lib .
Cree una nueva imagen a partir del contenedor. Dado que /opt/flink/lib se carga como classpath de forma predeterminada, ahora el archivo jar del controlador se encuentra en classpath.
Agregue maven-assembly-plugin al pom.xml de su proyecto maven. Vuelva a compilar su proyecto y obtenga un archivo jar con dependencias. En este contenedor, el controlador JDBC de PostgreSQL está empaquetado.