使用neo4j-spark-connector加载Dataframe遇类型不匹配及表不存在错误
先来说说你遇到的两个问题,其实第二个异常是根源,解决它之后第一个问题自然就消失了。
第二个异常:Map类型无法解析为String
你看到的java.util.Collections$UnmodifiableMap is not a valid external type for schema of string错误,是因为你的Cypher语句返回了整个properties(n)(一个Map对象),但Spark无法自动将这个Map类型正确映射到DataFrame的列类型上——connector错误地把它推断成了String类型,而实际返回的是UnmodifiableMap,类型不匹配导致编码失败。
解决方法:
方法1:展开属性为单独列(推荐)
不要返回整个properties(n),而是明确列出需要的属性,这样Spark能正确识别每个列的类型:
Dataset<Row> neo4jDF = neo4j.cypher("match (n:NodeType1) return n.id, n.otherProperty, ...", a).loadDataFrame(); neo4jDF.show(); // 现在应该能正常显示了
这样每个属性都是DataFrame的单独列,后续关联查询也能直接用neo4jtable.id匹配,非常方便。
方法2:显式指定Map类型的Schema
如果你确实需要保留整个properties作为Map列,可以手动定义Schema,告诉Spark这是一个Map类型:
import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.types.StructType; StructType schema = new StructType() .add("properties(n)", DataTypes.createMapType(DataTypes.StringType, DataTypes.StringType)); Dataset<Row> neo4jDF = neo4j.cypher("match (n:NodeType1) return properties(n)", a) .loadDataFrame(schema); neo4jDF.show();
这里假设你的properties都是字符串键值对,如果有其他类型,需要对应调整Map的value类型。
第一个异常:找不到neo4jtable视图
这个问题是因为之前的neo4jDF创建失败(第二个异常导致),所以createOrReplaceTempView("neo4jtable")实际上没有成功创建视图,Spark自然找不到它。当你解决了第二个异常,neo4jDF能正常创建后,再执行创建视图和关联查询就没问题了:
Dataset<Row> neo4jDF = neo4j.cypher("match (n:NodeType1) return n.id, n.otherProperty", a).loadDataFrame(); Dataset<Row> df2 = // 加载df2的代码 neo4jDF.createOrReplaceTempView("neo4jtable"); df2.createOrReplaceTempView("df2table"); Dataset<Row> joinedData = ss.sql("SELECT * from df2table JOIN neo4jtable ON df2table.id = neo4jtable.id"); joinedData.show();
另外建议用JOIN ... ON替代逗号分隔的写法,这是更规范的SQL关联语法。
内容的提问来源于stack exchange,提问作者Mahesha999

