如何在Hazelcast Jet JDBC中使用连接池?连接池配置序列化异常问题排查与解决方案
问题
我尝试在Hazelcast Jet JDBC中使用连接池,但无法实现该功能。我通过DataSource bean获取数据库连接,但无法正常工作。以下是我的代码:
Connection conn = ((DataSource)Ds.getBean("dataSourceName")).getConnection(); BatchSource<Object> jdbcSource = Sources .jdbc(() -> conn, (con, parallelism, index) -> { // query execution }, r -> this.mapResultSet1(r, metaData));
执行代码时出现如下错误:
java.lang.IllegalArgumentException: "newConnectionFn" must be serializable at com.hazelcast.jet.impl.util.Util.checkSerializable(Util.java:203) at com.hazelcast.jet.impl.connector.ReadJdbcP.supplier(ReadJdbcP.java:77) at com.hazelcast.jet.core.processor.SourceProcessors.readJdbcP(SourceProcessors.java:433) at com.hazelcast.jet.pipeline.Sources.jdbc(Sources.java:1327) at com.aivdata.impl.JDBCDataSource.readSource(JDBCDataSource.java:70) at com.aivdata.services.AivDataFactory.lambda$1(AivDataFactory.java:123) at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1384) at java.util.stream.ReferencePipeline$Head.forEach(ReferencePipeline.java:647) Caused by: java.io.NotSerializableException: org.apache.tomcat.dbcp.dbcp2.PoolingDataSource$PoolGuardConnectionWrapper at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1184) at java.io.ObjectOutputStream.writeArray(ObjectOutputStream.java:1378) at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1174) at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1548) at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1509) at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1432) at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1178) at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:348) at com.hazelcast.jet.impl.util.Util.checkSerializable(Util.java:201)
请问如何正确实现Hazelcast Jet与JDBC连接池的集成使用?
解决方案
这个问题的核心原因很清晰:你提前获取了Connection对象并把它放到了lambda里,而这个连接对象本身是不可序列化的——Hazelcast Jet需要把你的数据源逻辑序列化后分发到集群的各个节点执行,直接传递一个已经实例化的连接肯定会失败。
正确的做法是在lambda内部获取连接,而不是提前拿到连接再传入。另外,为了复用连接池的能力,你应该让DataSource本身是可序列化的(或者能在每个节点上获取到同一个数据源实例),然后在lambda里调用getConnection()。
具体步骤和代码示例:
- 确保你的
DataSourcebean是可序列化的,或者能在Jet集群的每个节点上通过Spring上下文获取到(比如配置成全局可用的bean)。 - 修改
jdbc()方法的第一个参数,让它在lambda内部获取连接:
BatchSource<Object> jdbcSource = Sources .jdbc( // 在lambda内部获取连接,这样每个节点都会从本地连接池拿连接 () -> ((DataSource)Ds.getBean("dataSourceName")).getConnection(), (con, parallelism, index) -> { // 在这里执行你的查询,比如: String sql = "SELECT * FROM your_table"; return con.prepareStatement(sql).executeQuery(); }, r -> this.mapResultSet1(r, metaData) );
额外注意事项:
- 如果你是在分布式集群环境中运行,要保证每个Jet节点都能访问到同一个数据库,并且连接池的配置(比如最大连接数、超时时间)在所有节点上保持一致。
- 不要忘记在查询完成后关闭ResultSet、Statement和Connection——不过Hazelcast Jet的JDBC源会自动帮你管理这些资源的关闭,所以你不用手动处理。
- 如果你的
DataSource本身不可序列化,你可以考虑把连接池的配置参数(比如URL、用户名、密码)序列化,然后在lambda内部重新初始化数据源:
// 假设你把连接配置封装成一个可序列化的对象 DbConfig dbConfig = new DbConfig(jdbcUrl, username, password); BatchSource<Object> jdbcSource = Sources .jdbc( () -> { // 每个节点上初始化连接池并获取连接 HikariDataSource ds = new HikariDataSource(); ds.setJdbcUrl(dbConfig.getJdbcUrl()); ds.setUsername(dbConfig.getUsername()); ds.setPassword(dbConfig.getPassword()); return ds.getConnection(); }, (con, parallelism, index) -> con.prepareStatement("SELECT * FROM your_table").executeQuery(), r -> this.mapResultSet1(r, metaData) );
这样修改后,你的lambda就变成了可序列化的,因为它不再持有一个不可序列化的Connection对象,而是在需要的时候才去获取连接,完美适配Hazelcast Jet的分布式执行模型。
内容的提问来源于stack exchange,提问作者user3458271
相关产品推荐
相关产品推荐

