【问题标题】:Can PySpark write array of strings to a database through a JDBC driver?PySpark 可以通过 JDBC 驱动程序将字符串数组写入数据库吗?
【发布时间】:2021-10-03 17:52:05
【问题描述】:

我正在使用 PySpark,我想将一个字符串数组插入到具有 JDBC 驱动程序的数据库中,但出现以下错误:

IllegalArgumentException: Can't get JDBC type for array<string>

当我有 UDF 的 ArrayType(StringType()) 格式时会发生此错误。当我尝试覆盖列类型时:

.option("createTableColumnTypes", "col1 ARRAY, col2 ARRAY, col3 ARRAY, col4 ARRAY")

我明白了:

DataType array is not supported.(line 1, pos 18)

这让我想知道问题是否在 Spark 3.1.2 中,其中没有数组映射,我必须将其转换为字符串,还是来自我正在使用的驱动程序?

作为参考,我使用CrateDB 作为数据库。这是它的驱动程序:crate.io/docs/jdbc/en/latest

【问题讨论】:

  • 你用的是什么驱动?
  • 刚刚更新了引用驱动的问题

标签: apache-spark jdbc pyspark apache-spark-sql cratedb


【解决方案1】:

可能切换到在 CrateDB 中使用 Postgres JDBC 而不是 crate-jdbc 可以解决您的问题。

使用 CrateCB 4.6.1 和 postgresql 42.2.23 测试的 PySpark 程序示例:

from pyspark.sql import Row

df = spark.createDataFrame([
    Row(a = [1, 2]),
    Row(a = [3, 4])
])
df

df.write \
  .format("jdbc") \
  .option("url", "jdbc:postgresql://<url-to-server>:5432/?sslmode=require") \
  .option("driver", "org.postgresql.Driver") \
  .option("dbtable", "<tableName>") \
  .option("user", "<username>") \
  .option("password", "<password>") \
  .save()

【讨论】:

  • 那行得通。现在我遇到了一些其他错误:在插入数组时,我只从 1,000,000 中得到前 1,000 行,它会抛出:Py4JJavaError: An error occurred while calling o536.save. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 3.0 failed 1 times, most recent failure: Lost task 0.0 in stage 3.0 (TID 3) (844a50665e58 executor driver): org.postgresql.util.PSQLException: ERROR: line 1:1: mismatched input 'ROLLBACK' expecting {'SELECT', 'DEALLOCATE', 'CREATE', 'ALTER', 'KILL', 'BEGIN'...
  • 当我尝试插入在 spark udf 中定义的浮点数时:Py4JJavaError: An error occurred while calling o1014.save. : org.postgresql.util.PSQLException: ERROR: Cannot find data type: float4
  • 由于某种原因,在定义与数据框行匹配的批量大小选项后,第一个问题能够得到解决。现在所有的行都被插入了
  • The Spark PostgresDialect 将 float 和 double 分别映射到 FLOAT4FLOAT8 (src),目前 CrateDB 不支持。不幸的是,不知道是否有解决方法
  • 刚刚尝试了“createTableColumnTypes”选项,它通过将其定义为浮点数来工作。感谢您的帮助!
【解决方案2】:

您能否尝试为数组添加数据类型,即ARRAY(TEXT)

.option("createTableColumnTypes", "col1 ARRAY(TEXT), col2 ARRAY(TEXT), col3 ARRAY(TEXT), col4 ARRAY(TEXT)")

SELECT ['Hello']::ARRAY;
--> SQLParseException[line 1:25: no viable alternative at input 'SELECT ['Hello']::ARRAY limit']
SELECT ['Hello']::ARRAY(TEXT);
--> SELECT OK, 1 record returned (0.002 seconds)

【讨论】:

  • 这样我得到:ParseException: mismatched input 'TEXT' expecting INTEGER_VALUE(line 1, pos 24)
猜你喜欢
  • 2021-10-14
  • 1970-01-01
  • 1970-01-01
  • 2021-02-09
  • 1970-01-01
  • 1970-01-01
  • 2017-08-25
  • 1970-01-01
  • 2020-10-09
相关资源
最近更新 更多