首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >将MySQL表转换为拼图时触发异常

将MySQL表转换为拼图时触发异常
EN

Stack Overflow用户
提问于 2016-10-27 17:08:17
回答 1查看 1.2K关注 0票数 5

我正在尝试使用spark 1.6.2将一个MySQL远程表转换为一个拼花文件。

该进程运行10分钟,填充内存,而不是从以下消息开始:

代码语言:javascript
复制
WARN NettyRpcEndpointRef: Error sending message [message = Heartbeat(driver,[Lscala.Tuple2;@dac44da,BlockManagerId(driver, localhost, 46158))] in 1 attempts
org.apache.spark.rpc.RpcTimeoutException: Futures timed out after [10 seconds]. This timeout is controlled by spark.executor.heartbeatInterval

在结束时,此错误失败:

代码语言:javascript
复制
ERROR ActorSystemImpl: Uncaught fatal error from thread [sparkDriverActorSystem-scheduler-1] shutting down ActorSystem [sparkDriverActorSystem]
java.lang.OutOfMemoryError: GC overhead limit exceeded

我使用以下命令在火花壳中运行它:

代码语言:javascript
复制
spark-shell --packages mysql:mysql-connector-java:5.1.26 org.slf4j:slf4j-simple:1.7.21 --driver-memory 12G

val dataframe_mysql = sqlContext.read.format("jdbc").option("url", "jdbc:mysql://.../table").option("driver", "com.mysql.jdbc.Driver").option("dbtable", "...").option("user", "...").option("password", "...").load()

dataframe_mysql.saveAsParquetFile("name.parquet")

我的最大执行器内存限制在12G。有没有办法强迫在“小”块中写入拼花文件,以释放内存?

EN

回答 1

Stack Overflow用户

回答已采纳

发布于 2016-10-28 10:04:57

问题似乎是,当您使用jdbc连接器读取数据时,没有定义分区。

默认情况下,从JDBC读取数据并不是分布式的,因此要启用分发版,必须设置手动分区。您需要一个列,它是一个很好的分区键,您必须预先了解分布情况。

很明显,这就是你的数据:

代码语言:javascript
复制
root 
|-- id: long (nullable = false) 
|-- order_year: string (nullable = false) 
|-- order_number: string (nullable = false) 
|-- row_number: integer (nullable = false) 
|-- product_code: string (nullable = false) 
|-- name: string (nullable = false) 
|-- quantity: integer (nullable = false) 
|-- price: double (nullable = false) 
|-- price_vat: double (nullable = false) 
|-- created_at: timestamp (nullable = true) 
|-- updated_at: timestamp (nullable = true)

在我看来,order_year是个很好的候选人。(根据你的评论,你似乎有20年了)

代码语言:javascript
复制
import org.apache.spark.sql.SQLContext

val sqlContext: SQLContext = ???

val driver: String = ???
val connectionUrl: String = ???
val query: String = ???
val userName: String = ???
val password: String = ???

// Manual partitioning
val partitionColumn: String = "order_year"

val options: Map[String, String] = Map("driver" -> driver,
  "url" -> connectionUrl,
  "dbtable" -> query,
  "user" -> userName,
  "password" -> password,
  "partitionColumn" -> partitionColumn,
  "lowerBound" -> "0",
  "upperBound" -> "3000",
  "numPartitions" -> "300"
)

val df = sqlContext.read.format("jdbc").options(options).load()

PS: partitionColumnlowerBoundupperBoundnumPartitions:如果指定了任何选项,则必须全部指定这些选项。

现在你可以把你的DataFrame保存到地板上了。

票数 6
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/40290478

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档