Tengo un archivo csv en Amazon s3 con un tamaño de 62 mb (114 000 filas). Lo estoy convirtiendo en un conjunto de datos de Spark y tomando las primeras 500 filas de él. El código es el siguiente;
DataFrameReader df = new DataFrameReader(spark).format("csv").option("header", true); Dataset<Row> set=df.load("s3n://"+this.accessId.replace("\"", "")+":"+this.accessToken.replace("\"", "")+"@"+this.bucketName.replace("\"", "")+"/"+this.filePath.replace("\"", "")+""); set.take(500)Toda la operación toma de 20 a 30 seg.
Ahora estoy intentando lo mismo, pero en lugar de usar csv, estoy usando la tabla mySQL con 119 000 filas. El servidor MySQL está en amazon ec2. El código es el siguiente;
String url ="jdbc:mysql://"+this.hostName+":3306/"+this.dataBaseName+"?user="+this.userName+"&password="+this.password; SparkSession spark=StartSpark.getSparkSession(); SQLContext sc = spark.sqlContext(); DataFrameReader df = new DataFrameReader(spark).format("csv").option("header", true); Dataset<Row> set = sc .read() .option("url", url) .option("dbtable", this.tableName) .option("driver","com.mysql.jdbc.Driver") .format("jdbc") .load(); set.take(500);Esto está tomando de 5 a 10 minutos. Estoy ejecutando chispa dentro de jvm. Utilizando la misma configuración en ambos casos.
Puedo usar la columna de partición, la partición numérica, etc., pero no tengo ninguna columna numérica y un problema más es que desconozco el esquema de la tabla.
Mi problema no es cómo disminuir el tiempo requerido, ya que sé que, en el caso ideal, Spark se ejecutará en el clúster, pero lo que no puedo entender es por qué esta gran diferencia de tiempo en los dos casos anteriores.
Este problema se ha tratado varias veces en StackOverflow:
y en fuentes externas:
así que solo para reiterar: de forma predeterminada, DataFrameReader.jdbc no distribuye datos ni lecturas. Utiliza hilo único, ejecutor único.
Para distribuir lee:
rangos de uso con lowerBound / upperBound :
Properties properties; Lower Dataset<Row> set = sc .read() .option("partitionColumn", "foo") .option("numPartitions", "3") .option("lowerBound", 0) .option("upperBound", 30) .option("url", url) .option("dbtable", this.tableName) .option("driver","com.mysql.jdbc.Driver") .format("jdbc") .load(); predicates
Properties properties; Dataset<Row> set = sc .read() .jdbc( url, this.tableName, {"foo < 10", "foo BETWWEN 10 and 20", "foo > 20"}, properties )1.Descargue una copia del conector JDBC para mysql. Creo que ya tienes uno.
wget http://central.maven.org/maven2/mysql/mysql-connector-java/5.1.38/mysql-connector-java-5.1.38.jar2. cree un archivo db-properties.flat en el siguiente formato
jdbcUrl=jdbc:mysql://${jdbcHostname}:${jdbcPort}/${jdbcDatabase} user=<username> password=<password>3. Primero cree una tabla vacía donde desee cargar los datos.
spark-shell --driver-class-path <your path to mysql jar> import java.io.{File, FileInputStream} import java.util.Properties import org.apache.spark.sql.SaveMode import org.apache.spark.sql.hive.HiveContext import org.apache.spark.{SparkConf, SparkContext} val sQLContext = new HiveContext(sc) import sQLContext.implicits._ import sQLContext.sql sQLContext.setConf("hive.exec.dynamic.partition", "true") sQLContext.setConf("hive.exec.dynamic.partition.mode", "nonstrict") val dbProperties = new Properties() dbProperties.load(new FileInputStream(new File("your_path_to/db- properties.flat"))) val jdbcurl = dbProperties.getProperty("jdbcUrl") val df1 = "(SELECT * FROM your_table_name) as s1" val df2 = sQLContext.read.jdbc(jdbcurl, df1, dbProperties) df2.write.format("orc").partitionBy("your_partition_column_name").mode(SaveMode.Append).saveAsTable("your_target_table_name")