【发布时间】:2019-04-09 13:54:24
【问题描述】:
我正在尝试从多个 s3 存储桶中读取文件。
最初存储桶将位于不同的区域,但看起来这是不可能的。
所以现在我已将另一个存储桶复制到与要读取的第一个存储桶相同的区域,也就是我正在执行 spark 作业的同一区域。
SparkSession 设置:
val sparkConf = new SparkConf()
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.registerKryoClasses(Array(classOf[Event]))
SparkSession.builder
.appName("Merge application")
.config(sparkConf)
.getOrCreate()
使用创建 SparkSession 中的 SQLContext 调用的函数:
private def parseEvents(bucketPath: String, service: String)(
implicit sqlContext: SQLContext
): Try[RDD[Event]] =
Try(
sqlContext.read
.option("codec", "org.apache.hadoop.io.compress.GzipCodec")
.json(bucketPath)
.toJSON
.rdd
.map(buildEvent(_, bucketPath, service).get)
)
主要流程:
for {
bucketOnePath <- buildBucketPath(config.bucketOne.name)
_ <- log(s"Reading events from $bucketOnePath")
bucketOneEvents: RDD[Event] <- parseEvents(bucketOnePath, config.service)
_ <- log(s"Enriching events from $bucketOnePath with originating region data")
bucketOneEventsWithRegion: RDD[Event] <- enrichEventsWithRegion(
bucketOneEvents,
config.bucketOne.region
)
bucketTwoPath <- buildBucketPath(config.bucketTwo.name)
_ <- log(s"Reading events from $bucketTwoPath")
bucketTwoEvents: RDD[Event] <- parseEvents(config.bucketTwo.name, config.service)
_ <- log(s"Enriching events from $bucketTwoPath with originating region data")
bucketTwoEventsWithRegion: RDD[Event] <- enrichEventsWithRegion(
bucketTwoEvents,
config.bucketTwo.region
)
_ <- log("Merging events")
mergedEvents: RDD[Event] <- merge(bucketOneEventsWithRegion, bucketTwoEventsWithRegion)
if mergedEvents.isEmpty() == false
_ <- log("Grouping merged events by partition key")
mergedEventsByPartitionKey: RDD[(EventsPartitionKey, Iterable[Event])] <- eventsByPartitionKey(
mergedEvents
)
_ <- log(s"Storing merged events to ${config.outputBucket.name}")
_ <- store(config.outputBucket.name, config.service, mergedEventsByPartitionKey)
} yield ()
我在日志中得到的错误(实际存储桶名称已更改但真实名称确实存在):
19/04/09 13:10:20 INFO SparkContext: Created broadcast 4 from rdd at MergeApp.scala:141
19/04/09 13:10:21 INFO FileSourceScanExec: Planning scan with bin packing, max size: 134217728 bytes, open cost is considered as scanning 4194304 bytes.
org.apache.spark.sql.AnalysisException: Path does not exist: hdfs:someBucket2
我的标准输出日志显示了主要代码在失败前走了多远:
Reading events from s3://someBucket/*/*/*/*/*.gz
Enriching events from s3://someBucket/*/*/*/*/*.gz with originating region data
Reading events from s3://someBucket2/*/*/*/*/*.gz
Merge failed: Path does not exist: hdfs://someBucket2
奇怪的是,无论我选择哪个存储桶,第一次读取总是有效的。 但无论存储桶如何,第二次读取总是失败。 这告诉我存储桶没有问题,但是在使用多个 s3 存储桶时会产生一些奇怪的现象。
我只能看到从单个 s3 存储桶读取多个文件的线程,而不是从多个 s3 存储桶读取多个文件的线程。
有什么想法吗?
【问题讨论】:
标签: apache-spark amazon-s3 amazon-emr