【发布时间】:2021-11-19 11:19:21
【问题描述】:
我创建了一个包含两个输入文件的 RDD,即 Edges 和 Node 文件。当我使用 Graph.fromEdge() 方法创建图表时,我得到了错误。有人可以帮我吗? inputEdgesTextFile 和 inputNodesTextFile 正在获取输入文本数据集。 在代码的最后一行,我收到错误。 我发布了我在代码中遇到的错误。
public static void main(String[] args) {
SparkConf conf = new SparkConf().setMaster("local").setAppName("GraphFileReadClass");
JavaSparkContext javaSparkContext = new JavaSparkContext(conf);
ClassTag<String> stringTag = scala.reflect.ClassTag$.MODULE$.apply(String.class);
ClassTag<String> intTag = scala.reflect.ClassTag$.MODULE$.apply(Integer.class);
$eq$colon$eq<String, String> tpEquals = scala.Predef.$eq$colon$eq$.MODULE$.tpEquals();
// Load an external Text File in Apache spark
//The text files number of lines and each line consists these structure
//SFEdge contains: | Edge_id integer | Source_Id integer | Destination_id integer | EdgeLength double |
//SFNodes contains: | Node_id integer | Longitude double | Latitude double |
JavaRDD<String> inputEdgesTextFile = javaSparkContext.textFile("./SFEdges.txt");
JavaRDD<String> inputNodesTextFile = javaSparkContext.textFile("./SFNodes.txt");
ArrayList<Tuple2<Integer, Integer>> nodes = new ArrayList<>();
ArrayList<Edge<Double>> edges = new ArrayList<>();
JavaRDD<NodesClass> nodesPart = inputNodesTextFile.mapPartitions(p -> {
ArrayList<NodesClass> nodeList = new ArrayList<NodesClass>();
int counter = 0;
while (p.hasNext()) {
String[] parts = p.next().split(" ");
NodesClass node = new NodesClass();
node.setNode_Id(Integer.parseInt(parts[0]));
node.setLongitude(Double.parseDouble(parts[1]));
node.setLatitude(Double.parseDouble(parts[2]));
nodes.add(new Tuple2<Integer, Integer>(counter, Integer.parseInt(parts[0])));
nodeList.add(node);
counter++;
}
return nodeList.iterator();
});
JavaRDD<Tuple2<Integer, Integer>> nodesRDD = javaSparkContext.parallelize(nodes);
nodesRDD.foreach(data -> System.out.print("Node details: " + data._1() + " " + data._2()));
JavaRDD<EdgeNetwork> edgesPart = inputEdgesTextFile.mapPartitions(p -> {
ArrayList<EdgeNetwork> edgeList = new ArrayList<EdgeNetwork>();
while (p.hasNext()) {
String[] parts = p.next().split(" ");
EdgeNetwork edgeNet = new EdgeNetwork();
edgeNet.setEdge_id(Integer.parseInt(parts[0]));
edgeNet.setSource_id(Integer.parseInt(parts[1]));
edgeNet.setDestination_id(Integer.parseInt(parts[2]));
edgeNet.setEdge_length(Double.parseDouble(parts[3]));
edges.add(new Edge<Double>(Long.parseLong(parts[1]), Long.parseLong(parts[2]),
Double.parseDouble(parts[3])));
edgeList.add(edgeNet);
}
return edgeList.iterator();
});
JavaRDD<Edge<Double>> edgesRDD = javaSparkContext.parallelize(edges);
Graph<String, Double> graph = Graph.fromEdges(edgesRDD.rdd(), " ", StorageLevel.MEMORY_ONLY(),
StorageLevel.MEMORY_ONLY(), stringTag, stringTag);
//The warning shows above this line for Graph<String,Double>
//Maybe the RDD that I have created has some errors. Please suggest me
【问题讨论】:
-
能否为您的输入文件添加示例以及您看到的错误消息?
-
@werner 我已将示例作为 cmets 添加到具有 txt 数据集结构的代码中。请检查一下。
标签: java apache-spark spark-graphx