【问题标题】:Dataproc conflict in hadoop temporary tableshadoop临时表中的Dataproc冲突
【发布时间】:2018-03-13 23:23:14
【问题描述】:

我有一个流程,可以在不同区域的 Dataproc 集群上并行执行 Spark 作业。为每个区域创建一个集群,执行 spark 作业并在完成后删除集群。

spark 作业使用 org.apache.spark.rdd.PairRDDFunctions.saveAsNewAPIHadoopDataset 方法传递 BigQuery Configuration 将数据保存在 BigQuery 表中。该作业将数据保存在多个表中,每个作业多次调用 saveAsNewAPIHadoopDataset 方法。

问题是,有时我会遇到由 Hadoop 临时 BigQuery 数据集中的冲突导致的错误,该数据集在内部创建以运行作业:

Exception in thread "main" com.google.api.client.googleapis.json.GoogleJsonResponseException: 409 Conflict
{
 "code" : 409,
 "errors" : [ {
   "domain" : "global",
   "message" : "Already Exists: Dataset <my-gcp-project>:<MY-DATASET>_hadoop_temporary_job_201802250620_0013",
   "reason" : "duplicate"
 } ],
 "message" : "Already Exists: Dataset <my-gcp-project>:<MY-DATASET>_hadoop_temporary_job_201802250620_0013"
}
    at com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:145)
    at com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)
    at com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)
    at com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)
    at com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1056)
    at com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)
    at com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)
    at com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)
    at com.google.cloud.hadoop.io.bigquery.BigQueryOutputCommitter.setupJob(BigQueryOutputCommitter.java:107)
    at org.apache.spark.rdd.PairRDDFunctions$$anonfun$saveAsNewAPIHadoopDataset$1.apply$mcV$sp(PairRDDFunctions.scala:1150)
    at org.apache.spark.rdd.PairRDDFunctions$$anonfun$saveAsNewAPIHadoopDataset$1.apply(PairRDDFunctions.scala:1078)
    at org.apache.spark.rdd.PairRDDFunctions$$anonfun$saveAsNewAPIHadoopDataset$1.apply(PairRDDFunctions.scala:1078)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
    at org.apache.spark.rdd.RDD.withScope(RDD.scala:358)
    at org.apache.spark.rdd.PairRDDFunctions.saveAsNewAPIHadoopDataset(PairRDDFunctions.scala:1078)
    at org.apache.spark.api.java.JavaPairRDD.saveAsNewAPIHadoopDataset(JavaPairRDD.scala:819)
    ...

上述异常中的时间戳 201802250620_0013_0013 后缀,我不确定它是否代表时间。

我的想法是,有时作业会同时运行并尝试创建名称中具有相同时间戳的数据集。在并行作业中或在另一个 saveAsNewAPIHadoopDataset 调用的同一作业中。

我们如何在不延迟作业执行的情况下避免此错误?

我使用的依赖是:

<dependency>
    <groupId>com.google.cloud.bigdataoss</groupId>
    <artifactId>bigquery-connector</artifactId>
    <version>0.10.2-hadoop2</version>
    <scope>provided</scope>
</dependency>

Dataproc 映像版本为 1.1

编辑 1:

我尝试使用 IndirectBigQueryOutputFormat,但现在我收到一条错误消息,指出 gcs 输出路径已经存在,即使在每个 saveAsNewAPIHadoopDataset 调用中传递不同的时间。

这是我的代码: SparkConf sc = new SparkConf().setAppName("MyApp");

try (JavaSparkContext jsc = new JavaSparkContext(sc)) {
    JavaPairRDD<String, String> filesJson = jsc.wholeTextFiles(jsonFolder, parts);
    JavaPairRDD<String, String> jsons = filesJson.flatMapToPair(new FileSplitter()).repartition(parts);
    JavaPairRDD<Object, JsonObject> objsJson = jsons.flatMapToPair(new JsonParser()).filter(t -> t._2() != null).cache();

    objsJson
    .filter(new FilterType(MSG_TYPE1))
    .saveAsNewAPIHadoopDataset(createConf("my-project:MY_DATASET.MY_TABLE1", "gs://my-bucket/tmp1"));

    objsJson
    .filter(new FilterType(MSG_TYPE2))
    .saveAsNewAPIHadoopDataset(createConf("my-project:MY_DATASET.MY_TABLE2", "gs://my-bucket/tmp2"));

    objsJson
    .filter(new FilterType(MSG_TYPE3))
    .saveAsNewAPIHadoopDataset(createConf("my-project:MY_DATASET.MY_TABLE3", "gs://my-bucket/tmp3"));

    // here goes another ingestion process. same code as above but diferrent params, parsers, etc.
}

Configuration createConf(String table, String outGCS) {
  Configuration conf = new Configuration();
  BigQueryOutputConfiguration.configure(conf, table, null, outGCS, BigQueryFileFormat.NEWLINE_DELIMITED_JSON, TextOutputFormat.class);
  conf.set("mapreduce.job.outputformat.class", IndirectBigQueryOutputFormat.class.getName());
  return conf;
}

【问题讨论】:

    标签: hadoop apache-spark google-cloud-dataproc


    【解决方案1】:

    我相信可能发生的情况是每个映射器都试图创建自己的数据集。这是相当低效的(并且消耗的每日配额与映射器的数量成正比)。

    另一种方法是使用IndirectBigQueryOutputFormat 作为输出类:

    IndirectBigQueryOutputFormat 的工作原理是首先将所有数据缓冲到 Cloud Storage 临时表中,然后在 commitJob 时通过一次操作将所有数据从 Cloud Storage 复制到 BigQuery。建议将其用于大型作业,因为它只需要每个 Hadoop/Spark 作业一个 BigQuery“加载”作业,而 BigQueryOutputFormat 为每个 Hadoop/Spark 任务执行一个 BigQuery 作业。

    在此处查看示例:https://cloud.google.com/dataproc/docs/tutorials/bigquery-connector-spark-example

    【讨论】:

    • 我尝试了类似示例,但调用了 saveAsNewAPIHadoopDataset 六次,每个目标表调用一次。对于每个调用,我都通过不同的表和outputGcsPath 参数传递不同的配置,但现在我收到一个错误消息,指出输出路径已经存在,甚至在每次调用时传递的路径都不同。
    • 我刚刚尝试了该示例,其中两个调用顺序写入两个不同的表/GCS 路径,它运行良好。你能分享你正在使用的代码吗?你又给BigQueryOutputConfiguration.configure打电话了吗?
    猜你喜欢
    • 2012-09-21
    • 1970-01-01
    • 2011-09-14
    • 1970-01-01
    • 1970-01-01
    • 2014-10-28
    • 2018-10-23
    • 1970-01-01
    • 2015-10-04
    相关资源
    最近更新 更多