使用PySpark通过JDBC读取Apache Ignite时遇Fetch size错误求助
PySpark读取Apache Ignite数据时Fetch Size报错的解决方法
问题场景
使用Python 3.7、PySpark 2.4.8,通过JDBC方式读取Apache Ignite中的country表,执行代码如下:
spark.read.format("jdbc").option("driver", "org.apache.ignite.IgniteJdbcThinDriver") .option("url", "jdbc:ignite:thin://172.19.0.1:10800;schema=fs_dev").option("dbtable", "country").load().show()
运行后触发如下错误:
java.sql.SQLException: Fetch size must be greater than zero. at org.apache.ignite.internal.jdbc.thin.JdbcThinStatement.setFetchSize(JdbcThinStatement.java:620) at org.apache.spark.sql.execution.datasources.jdbc.JDBCRDD.compute(JDBCRDD.scala:302) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346) at org.apache.spark.rdd.RDD.iterator(RDD.scala:310) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346) at org.apache.spark.rdd.RDD.iterator(RDD.scala:310) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346) at org.apache.spark.rdd.RDD.iterator(RDD.scala:310) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:346) at org.apache.spark.rdd.RDD.iterator(RDD.scala:310) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90) at org.apache.spark.scheduler.Task.run(Task.scala:123) at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:411) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:417) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:748)
解决方法
该错误源于Apache Ignite的JDBC Thin驱动强制要求fetch size必须大于0,而PySpark 2.x版本默认未设置有效fetch size值,导致驱动校验失败。只需在JDBC读取选项中显式添加fetchsize参数,设置一个大于0的数值即可,示例代码如下:
spark.read.format("jdbc") .option("driver", "org.apache.ignite.IgniteJdbcThinDriver") .option("url", "jdbc:ignite:thin://172.19.0.1:10800;schema=fs_dev") .option("dbtable", "country") .option("fetchsize", "1000") # 设置大于0的fetch size,可根据数据量调整 .load().show()
可以根据表数据量大小调整fetchsize的值,比如数据量较大时设置为5000或10000,平衡内存占用与读取效率。
内容的提问来源于stack exchange,提问作者Felix
相关产品推荐
相关产品推荐

