AbsaOSS / AbsaOSS/Jdbc2S

Error while using JDBC2S jar with PySpark

未關閉
#23 1 則留言 1 個 reaction 已指派 0 人 在 GitHub 檢視
主要語言
Scala
星號
10
分支
3
PR 合併指標
30 天內沒有已合併 PR

描述

Hi Jdbc2S team,
Thank you so much for this repo.

I'm working on a personal project where I need to fetch real-time data from the Database using Spark Structured Streaming API.
I tried 2 approaches listed below.

1. Downloaded the Jdbc2S jar and used it during SparkSession initiation.

```
format = "za.co.absa.spark.jdbc.streaming.source.providers.JDBCStreamingSourceProviderV1"
spark = SparkSession.builder.config("spark.jars","./pathtoJar1, ./jdbc2s_2.11-1.2.0.jar").appName("StreamingCount").getOrCreate()

```
When I try to execute the script, I get the following error.
```
Traceback (most recent call last):
File "SparkStreaming.py", line 14, in
count_sql = spark.readStream.format(format) \
File "/Users/subramaniav1/Library/Python/3.8/lib/python/site-packages/pyspark/sql/streaming.py", line 454, in load
return self._df(self._jreader.load())
File "/Users/subramaniav1/Library/Python/3.8/lib/python/site-packages/py4j/java_gateway.py", line 1321, in __call__
return_value = get_return_value(
File "/Users/subramaniav1/Library/Python/3.8/lib/python/site-packages/pyspark/sql/utils.py", line 111, in deco
return f(*a, **kw)
File "/Users/subramaniav1/Library/Python/3.8/lib/python/site-packages/py4j/protocol.py", line 326, in get_return_value
raise Py4JJavaError(
py4j.protocol.Py4JJavaError: An error occurred while calling o37.load.
: java.lang.NoClassDefFoundError: org/apache/spark/sql/sources/v2/DataSourceV2
```

I use Python 3.8.9, Scala 2.12.15 and Spark 3.2.1

2. I tried to follow the steps mentioned in https://stackoverflow.com/a/72158761 where someone was able to modify and use this with PySpark. I modified the object JDBCStreamingSourceV1 in JDBCStreamingSourceV1.scala as mentioned, and the pom.xml with updated scala version and tried to rebuild the jar using Maven.

I got build failure with the following error
```
Run starting. Expected test count is: 29
TestClassJDBCSingleOffset:
*** RUN ABORTED ***
java.lang.NoClassDefFoundError: org/apache/spark/sql/sources/v2/reader/streaming/Offset
at java.base/java.lang.ClassLoader.defineClass1(Native Method)
at java.base/java.lang.ClassLoader.defineClass(ClassLoader.java:1013)
at java.base/java.security.SecureClassLoader.defineClass(SecureClassLoader.java:150)
at java.base/jdk.internal.loader.BuiltinClassLoader.defineClass(BuiltinClassLoader.java:862)
at java.base/jdk.internal.loader.BuiltinClassLoader.findClassOnClassPathOrNull(BuiltinClassLoader.java:760)
at java.base/jdk.internal.loader.BuiltinClassLoader.loadClassOrNull(BuiltinClassLoader.java:681)
at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:639)
at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188)
at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:521)
at za.co.absa.spark.jdbc.streaming.source.offsets.TestClassJDBCSingleOffset.$anonfun$new$1(TestClassJDBCSingleOffset.scala:28)
...
Cause: java.lang.ClassNotFoundException: org.apache.spark.sql.sources.v2.reader.streaming.Offset
at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:641)
at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188)
at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:521)
at java.base/java.lang.ClassLoader.defineClass1(Native Method)
at java.base/java.lang.ClassLoader.defineClass(ClassLoader.java:1013)
at java.base/java.security.SecureClassLoader.defineClass(SecureClassLoader.java:150)
at java.base/jdk.internal.loader.BuiltinClassLoader.defineClass(BuiltinClassLoader.java:862)
at java.base/jdk.internal.loader.BuiltinClassLoader.findClassOnClassPathOrNull(BuiltinClassLoader.java:760)
at java.base/jdk.internal.loader.BuiltinClassLoader.loadClassOrNull(BuiltinClassLoader.java:681)
at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:639)
```

It'd be of great help if you can help with steps to use this repo with Pyspark

貢獻指南

這個儲存庫沒有索引到貢獻指南

評估

這個 Issue 還沒有評估資料。

把新 issue 寄到你的電子郵件信箱

精選適合新手參與的 GitHub issue 摘要。