Untitled

 avatar
unknown
json
10 months ago
7.8 kB
21
Indexable
/opt/conda/lib/python3.10/site-packages/pyspark/pandas/utils.py:1016: PandasAPIOnSparkAdviceWarning: If `index_col` is not specified for `to_spark`, the existing index is lost when converting to Spark DataFrame.
  warnings.warn(message, PandasAPIOnSparkAdviceWarning)
---------------------------------------------------------------------------
Py4JJavaError                             Traceback (most recent call last)
Cell In[37], line 87
     67 feature_combined = (
     68     feature_combined
     69     .join(
   (...)
     73     )
     74 )
     77 feature_combined = (
     78     feature_combined
     79     .select(
   (...)
     84     .withColumn("index", F.row_number().over(Window.orderBy("userID")))
     85 )
---> 87 DataFrameWriter.write_to_hdfs(
     88     project_name="credit_scoring",
     89     data_name="feature_a_score_v3_6_0",
     90     df_input=feature_combined,
     91     folder="risk",
     92     data_format="parquet",
     93     write_mode="overwrite",
     94     partition_level=PartitionLevel.columns,
     95     partition_columns=["p_month"],
     96 )

File /zda/zda-datatools/zda_datatools/datahub/decorator.py:249, in push_metadata.<locals>.outer_wrap.<locals>.wrap(*args, **kwargs)
    244             logger.error(
    245                 f"Exception during pushing metadata to DataHub and DQM [{dataset_id}]: {e}"
    246             )
    248 start_time = datetime.now()
--> 249 func(**kwargs)
    250 end_time = datetime.now()
    251 if track_writing:

File /zda/zda-datatools/zda_datatools/spark/load/dataframe_writer.py:479, in DataFrameWriter.write_to_hdfs(project_name, data_name, df_input, database_name, folder, data_format, write_mode, partition_level, partition_date, partition_columns, max_records_per_file, partition_count, repartition_mode, writing_options, skip_repartitioning, columnar_encryption_config)
    474     if not partition_columns:
    475         raise ValueError(
    476             "Partition data by columns but partition_columns is not provided. "
    477             f"partition_columns: {partition_columns}"
    478         )
--> 479     DataFrameWriter.__partition_by_column__(
    480         df, location, partition_columns, write_mode, data_format
    481     )
    482 else:
    483     options.update({"maxRecordsPerFile": max_records_per_file})

File /zda/zda-datatools/zda_datatools/spark/load/dataframe_writer.py:232, in DataFrameWriter.__partition_by_column__(df, location, partition_columns, write_mode, data_format)
    221 @staticmethod
    222 def __partition_by_column__(
    223     df: DataFrame,
   (...)
    227     data_format,
    228 ):
    229     logger.info(f"Partition data on HDFS by columns: {partition_columns}")
    230     df.write.mode(write_mode).format(data_format).partitionBy(
    231         partition_columns
--> 232     ).save(location)

File /opt/conda/lib/python3.10/site-packages/pyspark/sql/readwriter.py:1463, in DataFrameWriter.save(self, path, format, mode, partitionBy, **options)
   1461     self._jwrite.save()
   1462 else:
-> 1463     self._jwrite.save(path)

File /opt/conda/lib/python3.10/site-packages/py4j/java_gateway.py:1322, in JavaMember.__call__(self, *args)
   1316 command = proto.CALL_COMMAND_NAME +\
   1317     self.command_header +\
   1318     args_command +\
   1319     proto.END_COMMAND_PART
   1321 answer = self.gateway_client.send_command(command)
-> 1322 return_value = get_return_value(
   1323     answer, self.gateway_client, self.target_id, self.name)
   1325 for temp_arg in temp_args:
   1326     if hasattr(temp_arg, "_detach"):

File /opt/conda/lib/python3.10/site-packages/pyspark/errors/exceptions/captured.py:179, in capture_sql_exception.<locals>.deco(*a, **kw)
    177 def deco(*a: Any, **kw: Any) -> Any:
    178     try:
--> 179         return f(*a, **kw)
    180     except Py4JJavaError as e:
    181         converted = convert_exception(e.java_exception)

File /opt/conda/lib/python3.10/site-packages/py4j/protocol.py:326, in get_return_value(answer, gateway_client, target_id, name)
    324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client)
    325 if answer[1] == REFERENCE_TYPE:
--> 326     raise Py4JJavaError(
    327         "An error occurred while calling {0}{1}{2}.\n".
    328         format(target_id, ".", name), value)
    329 else:
    330     raise Py4JError(
    331         "An error occurred while calling {0}{1}{2}. Trace:\n{3}\n".
    332         format(target_id, ".", name, value))

Py4JJavaError: An error occurred while calling o5324.save.
: org.apache.spark.SparkException: Job aborted due to stage failure: ResultStage 108 (save at NativeMethodAccessorImpl.java:0) has failed the maximum allowable number of times: 4. Most recent failure reason:
org.apache.spark.shuffle.MetadataFetchFailedException: Missing an output location for shuffle 31 partition 0
	at org.apache.spark.MapOutputTracker$.validateStatus(MapOutputTracker.scala:1739)
	at org.apache.spark.MapOutputTracker$.$anonfun$convertMapStatuses$11(MapOutputTracker.scala:1686)
	at org.apache.spark.MapOutputTracker$.$anonfun$convertMapStatuses$11$adapted(MapOutputTracker.scala:1685)
	at scala.collection.Iterator.foreach(Iterator.scala:943)
	at scala.collection.Iterator.foreach$(Iterator.scala:943)
	at scala.collection.AbstractIterator.foreach(Iterator.scala:1431)
	at org.apache.spark.MapOutputTracker$.convertMapStatuses(MapOutputTracker.scala:1685)
	at org.apache.spark.MapOutputTrackerWorker.getMapSizesByExecutorIdImpl(MapOutputTracker.scala:1327)
	at org.apache.spark.MapOutputTrackerWorker.getMapSizesByExecutorId(MapOutputTracker.scala:1289)
	at org.apache.spark.shuffle.sort.SortShuffleManager.getReader(SortShuffleManager.scala:140)
	at org.apache.spark.shuffle.ShuffleManager.getReader(ShuffleManager.scala:63)
	at org.apache.spark.shuffle.ShuffleManager.getReader$(ShuffleManager.scala:57)
	at org.apache.spark.shuffle.sort.SortShuffleManager.getReader(SortShuffleManager.scala:73)
	at org.apache.spark.sql.execution.ShuffledRowRDD.compute(ShuffledRowRDD.scala:200)
	at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
	at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
	at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
	at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
	at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
	at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
	at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
	at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
	at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
	at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
	at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
	at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
	at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
	at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
	at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
	at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
	at org.apache.spark.scheduler.Task.run(Task.scala:141)
	at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
	at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
	at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
	at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
	at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
	at java.lang.Thread.run(Thread.java:748)
Editor is loading...
Leave a Comment