Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Persyaratan untuk protokol komit EMRFS S3-optimized
Protokol komit EMRFS S3-optimized digunakan ketika kondisi berikut terpenuhi:
-
Anda menjalankan tugas Spark yang menggunakan Spark, DataFrames, atau Datasets untuk mengganti tabel yang dipartisi.
-
Anda menjalankan pekerjaan Spark yang mode overwrite partisinya.
dynamic -
Multipart upload diaktifkan di Amazon EMR. Ini adalah opsi default. Untuk informasi selengkapnya, lihat Protokol komit EMRFS dan S3-optimized unggahan multibagian.
-
Cache sistem file untuk EMRFS diaktifkan. Ini adalah opsi default. Periksa apakah pengaturan
fs.s3.impl.disable.cachediatur kefalse. -
Dukungan sumber data bawaan Spark digunakan. Built-in dukungan sumber data digunakan dalam keadaan berikut:
-
Saat pekerjaan menulis ke sumber data bawaan atau tabel.
-
Saat pekerjaan menulis ke tabel Hive metastore Parquet. Ini terjadi ketika
spark.sql.hive.convertInsertingPartitionedTabledanspark.sql.hive.convertMetastoreParquetkeduanya diatur ke benar. Ini adalah pengaturan default. -
Saat pekerjaan menulis ke tabel ORC metastore Hive. Ini terjadi ketika
spark.sql.hive.convertInsertingPartitionedTabledanspark.sql.hive.convertMetastoreOrckeduanya diaturtrue. Ini adalah pengaturan default.
-
-
Operasi pekerjaan Spark yang menulis ke lokasi partisi default — misalnya,
${table_location}/k1=v1/k2=v2/— menggunakan protokol commit. Protokol tidak digunakan jika operasi pekerjaan menulis ke lokasi partisi kustom - misalnya, jika lokasi partisi kustom diatur menggunakanALTER TABLE SQLperintah. -
Nilai berikut untuk Spark mesti digunakan:
-
spark.sql.sources.commitProtocolClassharus diatur keorg.apache.spark.sql.execution.datasources.SQLEmrOptimizedCommitProtocol. Ini adalah pengaturan default untuk rilis Amazon EMR 5.30.0 dan yang lebih tinggi, dan 6.2.0 dan lebih tinggi. -
Opsi
partitionOverwriteModetulis atauspark.sql.sources.partitionOverwriteModeharus diatur kedynamic. Pengaturan default-nya adalahstatic.catatan
Parameter
partitionOverwriteModemenulis opsi diperkenalkan di Spark 2.4.0. Untuk Spark versi 2.3.2, disertakan dengan rilis Amazon EMR 5.19.0, mengaturspark.sql.sources.partitionOverwriteModeproperti. -
Jika pekerjaan Spark diganti ke tabel Hive metastore Parquet,,
spark.sql.hive.convertMetastoreParquetspark.sql.hive.convertInsertingPartitionedTable, danspark.sql.hive.convertMetastore.partitionOverwriteModeharus disetel ke.trueAda pengaturan default. -
Jika pekerjaan Spark diganti ke tabel ORC metastore Hive,,
spark.sql.hive.convertMetastoreOrcspark.sql.hive.convertInsertingPartitionedTable, danspark.sql.hive.convertMetastore.partitionOverwriteModeharus disetel ke.trueAda pengaturan default.
-
contoh- Mode penggantian partisi dinamis
Dalam contoh Scala ini, pengoptimalan dipicu. Pertama, Anda mengatur partitionOverwriteMode properti kedynamic. Ini hanya mengganti partisi tempat Anda menulis data. Kemudian, Anda menentukan kolom partisi dinamis dengan partitionBy dan mengatur mode tulis keoverwrite.
val dataset = spark.range(0, 10) .withColumn("dt", expr("date_sub(current_date(), id)")) dataset.write.mode("overwrite") // "overwrite" instead of "insert" .option("partitionOverwriteMode", "dynamic") // "dynamic" instead of "static" .partitionBy("dt") // partitioned data instead of unpartitioned data .parquet("s3://amzn-s3-demo-bucket1/output") // "s3://" to use Amazon EMR file system, instead of "s3a://" or "hdfs://"
Ketika protokol komit EMR S3-optimized FS tidak digunakan
Umumnya, protokol komit EMR S3-optimized FS bekerja sama dengan protokol komit Spark default open source,. org.apache.spark.sql.execution.datasources.SQLHadoopMapReduceCommitProtocol Optimasi tidak akan terjadi dalam situasi berikut.
| Situasi | Mengapa protokol commit tidak digunakan |
|---|---|
| Saat Anda menulis ke HDFS | Protokol commit hanya mendukung penulisan ke Amazon S3 menggunakan EMRFS. |
| Saat Anda menggunakan sistem file S3A | Protokol commit hanya mendukung EMRFS. |
| Saat Anda menggunakan MapReduce atau API RDD Spark | Protokol komit hanya mendukung penggunaan SparkSQL, DataFrame, atau Dataset API. |
| Ketika penggantian partisi dinamis tidak dipicu | Protokol commit hanya mengoptimalkan kasus penggantian partisi dinamis. Untuk kasus lain, lihatGunakan komitter EMRFS S3-optimized. |
Contoh Scala berikut menunjukkan beberapa situasi tambahan yang menjadi delegasi protokol EMR S3-optimized FS. SQLHadoopMapReduceCommitProtocol
contoh- Mode penggantian partisi dinamis dengan lokasi partisi khusus
Dalam contoh ini, program Scala mengganti dua partisi dalam mode overwrite partisi dinamis. Satu partisi memiliki lokasi partisi kustom. Partisi lain menggunakan lokasi partisi default. Protokol komit EM S3-optimized RFS hanya meningkatkan partisi yang menggunakan lokasi partisi default.
val table = "dataset" val inputView = "tempView" val location = "s3://bucket/table" spark.sql(s""" CREATE TABLE $table (id bigint, dt date) USING PARQUET PARTITIONED BY (dt) LOCATION '$location' """) // Add a partition using a custom location val customPartitionLocation = "s3://bucket/custom" spark.sql(s""" ALTER TABLE $table ADD PARTITION (dt='2019-01-28') LOCATION '$customPartitionLocation' """) // Add another partition using default location spark.sql(s"ALTER TABLE $table ADD PARTITION (dt='2019-01-29')") def asDate(text: String) = lit(text).cast("date") spark.range(0, 10) .withColumn("dt", when($"id" > 4, asDate("2019-01-28")).otherwise(asDate("2019-01-29"))) .createTempView(inputView) // Set partition overwrite mode to 'dynamic' spark.sql(s"SET spark.sql.sources.partitionOverwriteMode=dynamic") spark.sql(s"INSERT OVERWRITE TABLE $table SELECT * FROM $inputView")
Kode Scala menciptakan objek Amazon S3 berikut:
custom/part-00001-035a2a9c-4a09-4917-8819-e77134342402.c000.snappy.parquet custom_$folder$ table/_SUCCESS table/dt=2019-01-29/part-00000-035a2a9c-4a09-4917-8819-e77134342402.c000.snappy.parquet table/dt=2019-01-29_$folder$ table_$folder$
catatan
Menulis ke lokasi partisi khusus di versi Spark sebelumnya dapat mengakibatkan hilangnya data. Dalam contoh ini, partisi dt='2019-01-28' akan hilang. Untuk detail selengkapnya, lihat SPARK-35106
Ketika menulis ke partisi di lokasi kustom, Spark menggunakan algoritma komit mirip dengan contoh sebelumnya, yang diuraikan di bawah ini. Seperti contoh sebelumnya, algoritma menghasilkan penggantian nama berurutan, yang dapat berdampak negatif pada kinerja.
Algoritma di Spark 2.4.0 mengikuti langkah-langkah berikut:
-
Ketika menulis output ke partisi di lokasi kustom, tugas menulis ke file di bawah direktori pementasan Spark ini, yang dibuat di bawah lokasi output akhir. Nama file termasuk UUID acak untuk melindungi terhadap tabrakan file. Upaya tugas melacak setiap file bersama dengan path output akhir yang diinginkan.
-
Ketika tugas selesai berhasil, menyediakan driver dengan file dan akhir yang diinginkan output jalan mereka.
-
Setelah semua tugas selesai, pekerjaan commit fase berurutan mengganti nama semua file yang ditulis untuk partisi di lokasi kustom ke jalur output akhir mereka.
-
Direktori pementasan dihapus sebelum pekerjaan komit fase selesai.