Spark Structured Streaming写入Neo4j报不支持流写入问题咨询
问题根因
抛出java.lang.UnsupportedOperationException: Data source org.neo4j.spark.DataSource does not support streamed writing错误的核心原因,是使用了错误的Neo4j数据源入口类,附带代码中还有3处笔误/配置错误会导致后续运行异常:
- 核心错误:4.1.2版本的Neo4j Connector for Apache Spark中,批处理和流处理的数据源实现是分离的,你写入时指定的
org.neo4j.spark.DataSource是批处理专用实现,没有实现Spark Structured Streaming要求的流写入接口,因此会抛出不支持流写入的异常。 - 配置拼写错误:读取MongoDB流时,读偏好配置的key缺了首字母
s,写成了park.mongodb.read.readPreference.name,该配置不会生效。 - 变量名不匹配:读取MongoDB变更流得到的DataFrame命名为
dfTxn,写入时调用的是未定义的dfPaymentTx变量,即使数据源类配置正确也会抛出变量不存在的错误。 - 无效配置项:流写入场景下
save.mode是无效配置,Neo4j连接器的节点/关系写入策略需要通过node.save.mode、relationship.save.mode单独指定,流场景默认的Append语义已经匹配你的需求,不需要额外配置通用save模式。
修复方案
按以下步骤调整代码即可正常运行:
- 替换流写入的数据源类,将
format的参数值从批处理专用的org.neo4j.spark.DataSource,替换为流处理专用的org.neo4j.spark.streaming.Neo4jStreamingDataSource - 修正MongoDB读偏好配置的拼写,补全key开头的
s - 统一DataFrame变量名,写入时使用实际加载得到的
dfTxn - 删除无效的
save.mode配置,如果需要调整节点写入的冲突策略,可显式添加node.save.mode配置 - 启动任务前清空旧的checkpoint目录,避免之前失败任务的残留状态干扰新任务运行
修复后完整代码示例
启动spark-shell命令
spark-shell --packages org.mongodb.spark:mongo-spark-connector:10.0.2,org.neo4j:neo4j-connector-apache-spark_2.12:4.1.2_for_spark_3
读取MongoDB变更流代码
val dfTxn = spark.readStream.format("mongodb") .option("spark.mongodb.connection.uri", "mongodb://<IP>:<PORT>") .option("spark.mongodb.database", "test") .option("spark.mongodb.collection", "txn") .option("spark.mongodb.read.readPreference.name","primaryPreferred") .option("spark.mongodb.change.stream.publish.full.document.only", "true") .option("forceDeleteTempCheckpointLocation", "true") .load()
写入Neo4j流代码
val query = dfTxn.writeStream .format("org.neo4j.spark.streaming.Neo4jStreamingDataSource") .option("url", "bolt://<IP>:<PORT>") .option("checkpointLocation", "/tmp/checkpoint/myCheckPoint") .option("labels", "Account") .option("node.keys", "txn_snd") .start()
额外注意事项
4.1.2版本的Neo4j连接器仅支持Spark 3.0~3.2版本,如果你的Spark版本是3.3及以上,会出现类兼容问题,即使配置了正确的流数据源类也可能报错,这种情况需要升级Neo4j连接器到对应适配Spark版本的发行版。
内容的提问来源于stack exchange,提问作者abaschkim
相关产品推荐
相关产品推荐

