Business
Jobs
  • About Us
  • Solutions
    • Job Postings
      Post your job and receive qualified candidates in 48h.
    • Candidate Assessments
      500+ technical and psychological tests, plus anti-fraud.
    • Headhunting
      Tailor-made executive search from start to finish.
    • Payroll + EOR
      Payroll dispersal and EOR across 15+ LATAM countries.
  • Pricing
  • Jobs

0

79
Views
La conversión de una tabla mysql a un conjunto de datos de chispa es muy lenta en comparación con el mismo archivo csv

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.

over 4 years ago · Santiago Trujillo
2 answers
Answer question

0

Este problema se ha tratado varias veces en StackOverflow:

  • ¿Cómo mejorar el rendimiento de los trabajos lentos de Spark usando DataFrame y la conexión JDBC?
  • chispa jdbc df límite... ¿qué está haciendo?
  • ¿Cómo usar la fuente JDBC para escribir y leer datos en (Py) Spark?

y en fuentes externas:

  • https://github.com/awesome-spark/spark-gotchas/blob/master/05_spark_sql_and_dataset_api.md#parallelizing-reads

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 )
over 4 years ago · Santiago Trujillo Report

0

Por favor, siga los pasos a continuación

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.jar

2. 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.

invocar spark shell con clase de controlador

 spark-shell --driver-class-path <your path to mysql jar>

luego importe todo el paquete requerido

 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}

iniciar un contexto de colmena o un contexto sql

 val sQLContext = new HiveContext(sc) import sQLContext.implicits._ import sQLContext.sql

establecer algunas de las propiedades

 sQLContext.setConf("hive.exec.dynamic.partition", "true") sQLContext.setConf("hive.exec.dynamic.partition.mode", "nonstrict")

Cargue las propiedades de mysql db desde el archivo

 val dbProperties = new Properties() dbProperties.load(new FileInputStream(new File("your_path_to/db- properties.flat"))) val jdbcurl = dbProperties.getProperty("jdbcUrl")

cree una consulta para leer los datos de su tabla y pásela al método de lectura de #sqlcontext. aquí es donde puede administrar su cláusula where

 val df1 = "(SELECT * FROM your_table_name) as s1"

pase jdbcurl, seleccione la consulta y las propiedades de db para leer el método

 val df2 = sQLContext.read.jdbc(jdbcurl, df1, dbProperties)

escríbelo en tu mesa

 df2.write.format("orc").partitionBy("your_partition_column_name").mode(SaveMode.Append).saveAsTable("your_target_table_name")
over 4 years ago · Santiago Trujillo Report
Answer question
Find remote jobs

Discover the new way to find a job!

Top jobs
Top job categories
Business
Post vacancy Pricing Sales
Legal
Terms and conditions Privacy policy
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Show me some job opportunities
There's an error!