【发布时间】:2020-04-07 11:41:11
【问题描述】:
我有一个周期性的 spark-scala 任务,用于将数据从 Hive 传输到 MySQL。
table的结构可以简单的看成:
+------+------+
| id | name |
+------+------+
然后因为hive表太大,所以只好共享mysql表。
所以这是我目前的解决方案:
- 准备 MySQL 表:
mysql> show tables;
+-------------------+
| Tables_in_test_db |
+-------------------+
| shared_0 |
| shared_1 |
| shared_2 |
| shared_3 |
| shared_4 |
| shared_5 |
+-------------------+
- 从 Hive 加载数据并进行一些转换操作,然后生成我想要的数据帧
val data = List((0, "a"), (11, "b"), (22, "c"), (33, "d"), (44, "e"))
val total = spark.sparkContext.parallelize(data)
.toDF("id", "name")
.withColumn("hashCode", hash($"id")%5)
- 根据
hashCode列将数据保存到MySQL表中
(0 to 5).foreach(hashCode => {
val df = total.where($"hashCode" === hashCode).select("id", "name")
df.write
.mode(SaveMode.Append)
.jdbc(jdbcUrl, s"shared_$hashCode", connectionProperties)
})
这很好用,但我是 spark 的新手,所以我想知道有没有更好的方法来实现我想要的??
更新:
这是我的完整代码:
val jdbcHostname = "localhost"
val jdbcPort = 3306
val jdbcDatabase = "test_db"
val jdbcUrl = s"jdbc:mysql://${jdbcHostname}:${jdbcPort}/${jdbcDatabase}"
val connectionProperties = new Properties()
connectionProperties.put("user", "user")
connectionProperties.put("password", "password")
val spark = SparkSession.builder().master("local[*]").appName("test").getOrCreate()
import spark.implicits._
val data = List((0, "a"), (11, "b"), (22, "c"), (33, "d"), (44, "e"))
val total = spark.sparkContext.parallelize(data)
.toDF("id", "name")
.withColumn("hashCode", hash($"id")%5)
(0 to 5).foreach(hashCode => {
val df = total.where($"hashCode" === hashCode).select("id", "name")
df.write
.mode(SaveMode.Append)
.jdbc(jdbcUrl, s"shared_$hashCode", connectionProperties)
})
【问题讨论】:
标签: apache-spark apache-spark-sql