Scala数组Distinct返回空值及Spark GraphX代码运行问题咨询
Hey there! Let's break down solutions for your two Scala/Spark problems:
distinct If your array's distinct call is returning empty strings, here are the most common fixes:
First, verify your original array content
Print out all elements to check if you're starting with mostly empty strings or hidden whitespace:val yourArray = Array("", " ", "\t", "") println(yourArray.mkString("|")) // Shows you exactly what's in the arrayClean hidden whitespace before deduplication
Often, "empty" looking strings are actually filled with spaces/tabs. Trim them first:val cleanedArray = yourArray.map(_.trim).distinct // This turns Array("", " ", "\t", "") into Array("")Check for non-printable characters
If trimming doesn't help, use regex to strip all non-visible characters:val sanitizedArray = yourArray.map(_.replaceAll("\\p{C}", "")).distinctValidate your data source
If after cleaning you still get empty strings, your source data might actually be returning empty values. Double-check where the array is coming from (e.g., CSV reads, API responses).
Looking at your code snippet, here are the key issues and fixes:
First: Define the Flight Case Class
You're trying to return a Flight object, but you haven't defined the class yet. In Spark Shell, define this before your parseFlight function:
case class Flight( origin: String, dest: String, carrier: String, tailNum: String, depDelay: Int, depTime: Long, arrTime: String, arrDelay: Long, cancelled: String, distance: Double, taxiOut: Double, taxiIn: Double, depDelayMinutes: Double, arrDelayMinutes: Double, airTime: Double, flightNum: Double, year: Int ) // Adjust field names to match your actual data meaning!
Second: Add Error Handling to parseFlight
CSV data often has malformed lines or missing values. Wrap your parsing in a try-catch to avoid crashes:
def parseFlight(str: String): Option[Flight] = { try { val line = str.split(",") // Make sure we have at least 17 columns (indices 0-16) if (line.length >= 17) { Some(Flight( line(0), line(1), line(2), line(3), line(4).toInt, line(5).toLong, line(6), line(7).toLong, line(8), line(9).toDouble, line(10).toDouble, line(11).toDouble, line(12).toDouble, line(13).toDouble, line(14).toDouble, line(15).toDouble, line(16).toInt )) } else { println(s"Skipping short line: $str") None } } catch { case e: NumberFormatException => println(s"Invalid number in line: $str | Error: ${e.getMessage}") None case e: Exception => println(s"Failed to parse line: $str | Error: ${e.getMessage}") None } }
Third: Complete Your RDD Pipeline
Your code cuts off at val flig...—finish it by mapping and filtering valid flights:
val textRDD = sc.textFile("/user/user01/data/rita2014jan.csv") // Skip header if your CSV has one! val flightsRDD = textRDD .filter(!_.startsWith("origin")) // Adjust header text to match your CSV .map(parseFlight) .filter(_.isDefined) .map(_.get)
Fourth: Verify File Access
Make sure the HDFS path exists and you have permissions:
// Test if the file is readable textRDD.count() // If this throws an error, check path/permissions
Bonus: Start Building Your GraphX Graph
Once you have valid flights, you can create vertices and edges for GraphX:
import org.apache.spark.graphx._ // Create airport vertices (ID, airport code) val airportVertices = flightsRDD .flatMap(f => List(f.origin, f.dest)) .distinct() .map(code => (code.hashCode.toLong, code)) // Create flight edges (origin ID -> dest ID, flight count) val flightEdges = flightsRDD .map(f => Edge(f.origin.hashCode.toLong, f.dest.hashCode.toLong, 1L)) // Build the graph val flightGraph = Graph(airportVertices, flightEdges)
内容的提问来源于stack exchange,提问作者vikash

