View a markdown version of this page

Menjalankan PySpark pekerjaan pada tabel yang dikonfigurasi menggunakan template PySpark analisis - AWS Clean Rooms

Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.

Menjalankan PySpark pekerjaan pada tabel yang dikonfigurasi menggunakan template PySpark analisis

Prosedur ini menunjukkan cara menggunakan template PySpark analisis di AWS Clean Rooms konsol untuk menganalisis tabel yang dikonfigurasi dengan aturan analisis kustom.

Untuk menjalankan PySpark pekerjaan pada tabel yang dikonfigurasi menggunakan template PySpark analisis

Masuk ke Konsol Manajemen AWS dan buka AWS Clean Rooms konsol di https://console.aws.amazon.com/cleanrooms.

  1. Di panel navigasi kiri, pilih Kolabor asi.

  2. Pilih kolaborasi yang memiliki status kemampuan anggota Anda dari J alankan pekerjaan.

  3. Pada tab Analisis, di bawah bagian Tabel, lihat tabel dan jenis aturan analisis terkait (A turan analisis khusus).

    catatan

    Jika Anda tidak melihat tabel yang Anda harapkan dalam daftar, itu mungkin karena alasan berikut:

  4. Di bawah bagian Analisis, untuk mode Analisis, pilih J alankan template analisis.

  5. Pilih template PySpark analisis dari daftar dropdown template Analisis.

    Parameter dari template PySpark analisis akan secara otomatis terisi dalam Definisi.

  6. Jika template analisis memiliki parameter yang ditentukan, di bawah Param eter, berikan nilai untuk parameter:

    1. Untuk setiap parameter, lihat nama Parameter dan nilai Default (jika dikonfigurasi).

    2. Masukkan Nilai untuk setiap parameter yang ingin Anda ganti.

      catatan

      Jika Anda tidak memberikan nilai tetapi nilai default ada, nilai default akan digunakan.

    penting

    Nilai parameter dapat hingga 1.000 karakter dan mendukung UTF-8 pengkodean. Semua nilai parameter diperlakukan sebagai string dan diteruskan ke skrip pengguna Anda melalui objek konteks.

    Pastikan skrip pengguna Anda memvalidasi dan menangani nilai parameter dengan aman. Untuk informasi selengkapnya tentang penanganan parameter aman, lihatBekerja dengan parameter dalam templat PySpark analisis.

  7. Tentukan jenis Pekerja yang didukung dan Jumlah pekerja.

    Gunakan tabel berikut untuk menentukan jenis dan jumlah atau pekerja yang Anda butuhkan untuk kasus penggunaan Anda.

    Jenis pekerja vCPU Memori (GB) Penyimpanan (GB) Jumlah pekerja Total Unit Pemrosesan Kamar Bersih (CRPU)
    CR.1X (default) 4 30 100 4 8
    128 256
    CR.4X 16 120 400 4 32
    32 256
    catatan

    Jenis pekerja yang berbeda dan jumlah pekerja memiliki biaya terkait. Untuk mempelajari lebih lanjut tentang harga, lihat AWS Clean Rooms harga.

  8. Tentukan properti Spark yang didukung.

    1. Pilih Tambahkan properti Spark.

    2. Pada kotak dialog properti Spark, pilih nama Pro perti dari daftar dropdown dan masukkan Nilai.

    Tabel berikut memberikan definisi untuk setiap properti.

    Untuk informasi selengkapnya tentang properti Spark, lihat Spark Properties dalam dokumentasi Apache Spark.

    catatan

    Anda dapat mengonfigurasi maksimal 50 properti Spark. Setiap nilai properti dapat mencapai 500 karakter.

    Nama Properti Deskripsi nilai default

    Spark.task.maxFailures

    Mengontrol berapa kali berturut-turut tugas dapat gagal sebelum pekerjaan gagal. Membutuhkan nilai yang lebih besar dari atau sama dengan 1. Jumlah percobaan ulang yang diizinkan sama dengan nilai ini dikurangi 1. Jumlah kegagalan diatur ulang jika ada upaya yang berhasil. Kegagalan di berbagai tugas tidak menumpuk menuju batas ini.

    4

    spark.sql.files.max PartitionBytes

    Menetapkan jumlah maksimum byte untuk dikemas ke dalam satu partisi saat membaca dari sumber berbasis file seperti Parquet, JSON, dan ORC.

    128MB

    spark.hadoop.fs.s3.maxRetries

    Menetapkan jumlah maksimum upaya coba ulang untuk operasi file Amazon S3.

    (tidak ada)

    spark.network.timeout

    Menetapkan batas waktu default untuk semua interaksi jaringan. Mengganti setelan batas waktu berikut jika tidak dikonfigurasi:

    • spark.storage.blok ManagerHeartbeatTimeoutMs

    • Spark.shuffle.io.ConnectionTimeout

    • Spark.rpc.AskTimeout

    • Spark.RPC.LookupTimeout

    120-an

    spark.rdd.kompres

    Menentukan apakah akan mengompres partisi RDD serial menggunakan spark.io.compression.codec. Berlaku untuk StorageLevel.MEMORY _ONLY_SER di Java dan Scala, atau StorageLevel.MEMORY _ONLY di Python. Mengurangi ruang penyimpanan tetapi membutuhkan waktu pemrosesan CPU tambahan.

    false

    spark.shuffle.spill.kompres

    Menentukan apakah akan mengompres data shuffle spill menggunakan spark.io.compression.codec.

    true

    spark.shuffle.kompres

    Menentukan apakah akan mengompres file output peta. Kompresi menggunakan spark.io.compression.codec.

    true

    spark.shuffle.service.index.cache.size

    Menetapkan batas ukuran cache, dalam byte kecuali ditentukan lain.

    100m

    Spark.shuffle.io.maxRetries

    Menetapkan jumlah maksimum percobaan ulang untuk pengambilan yang gagal karena pengecualian. IO-related

    3

    spark.shuffle.io.RetryWait

    Menetapkan waktu tunggu antara percobaan ulang pengambilan. Penundaan maksimum yang disebabkan oleh percobaan ulang adalah 15 detik secara default, dihitung sebagai maxRetries * retryWait.

    5 detik

    Spark.shuffle.io.ConnectionTimeout

    Menetapkan batas waktu untuk koneksi yang dibuat antara server shuffle dan klien untuk ditandai sebagai idle dan ditutup jika masih ada permintaan pengambilan yang belum dibayar tetapi tidak ada lalu lintas di saluran.

    (nilai spark.network.timeout)

    spark.driver.max ResultSize

    Menetapkan batas ukuran total hasil serial dari semua partisi untuk setiap tindakan Spark, dalam byte. Setidaknya harus 1M, atau 0 untuk tidak terbatas.

    1g

    spark.memory.fraksi

    Menetapkan fraksi (ruang tumpukan - 300MB) yang digunakan untuk eksekusi dan penyimpanan. Semakin rendah nilai ini, semakin sering terjadi tumpahan dan penggusuran data yang di-cache. Dianjurkan untuk meninggalkan ini pada nilai default.

    0.6

    spark.scheduler.mode

    Menetapkan mode penjadwalan antara pekerjaan yang dikirimkan ke yang sama SparkContext. Dapat diatur ke FAIR untuk menggunakan pembagian yang adil alih-alih mengantri pekerjaan satu demi satu. Nilai yang didukung: FAIR, FIFO.

    FIFO

    spark.sql.adaptive.advisory PartitionSizeInBytes

    Menetapkan ukuran target dalam byte untuk partisi shuffle selama pengoptimalan adaptif saat spark.sql.adaptive.enabled benar. Mengontrol ukuran partisi saat menggabungkan partisi kecil atau memisahkan partisi miring.

    (nilai spark. PostShuffleInputSize sql.adaptive.shuffle.target)

    spark.sql.adaptive.otomatis BroadcastJoinThreshold

    Menetapkan ukuran tabel maksimum dalam byte untuk penyiaran ke node pekerja selama bergabung. Hanya berlaku dalam kerangka adaptif. Menggunakan nilai default yang sama dengan spark. BroadcastJoinThreshold sql.auto. Setel ke -1 untuk menonaktifkan penyiaran.

    (tidak ada)

    spark.sql.adaptive.coalesce Partitions.enabled

    Menentukan apakah akan menggabungkan partisi shuffle yang berdekatan berdasarkan spark.sql.adaptive.advisory untuk mengoptimalkan ukuran tugas. PartitionSizeInBytes Membutuhkan spark.sql.adaptive.enabled agar benar.

    true

    spark.sql.adaptive.coalesce Partitions.initialPartitionNum

    Mendefinisikan jumlah awal partisi shuffle sebelum penggabungan. Memerlukan spark.sql.adaptive.enabled dan spark.sql.adaptive.coalesce agar benar. Partitions.enabled Defaultnya adalah nilai spark.sql.shuffle.partition.

    (tidak ada)

    spark.sql.adaptive.coalesce Partitions.minPartitionSize

    Menetapkan ukuran minimum untuk partisi shuffle gabungan untuk mencegah partisi menjadi terlalu kecil selama pengoptimalan adaptif.

    1 MB

    spark.sql.adaptive.coalesce Partitions.parallelismFirst

    Menentukan apakah akan menghitung ukuran partisi berdasarkan paralelisme cluster alih-alih spark.sq PartitionSizeInBytes l.adaptive.advisory selama penggabungan partisi. Menghasilkan ukuran partisi yang lebih kecil dari ukuran target yang dikonfigurasi untuk memaksimalkan paralelisme. Sebaiknya setel ini ke false pada cluster sibuk untuk meningkatkan pemanfaatan sumber daya dengan mencegah tugas kecil yang berlebihan.

    true

    spark.sql.adaptive.enabled

    Menentukan apakah akan mengaktifkan eksekusi kueri adaptif untuk mengoptimalkan kembali rencana kueri selama eksekusi kueri, berdasarkan statistik runtime yang akurat.

    true

    spark.sql.adaptive.force OptimizeSkewedJoin

    Menentukan apakah akan memaksa mengaktifkan OptimizeSkewedJoin bahkan jika itu memperkenalkan shuffle tambahan.

    false

    spark.sql.adaptive.lokal ShuffleReader.enabled

    Menentukan apakah akan menggunakan pembaca shuffle lokal saat partisi shuffle tidak diperlukan, seperti setelah mengonversi dari gabungan sort-merge ke gabungan broadcast-hash. Membutuhkan spark.sql.adaptive.enabled agar benar.

    true

    spark.sql.adaptive.max ShuffledHashJoinLocalMapThreshold

    Menetapkan ukuran partisi maksimum dalam byte untuk membangun peta hash lokal. Memprioritaskan gabungan hash yang diacak daripada gabungan sort-merge ketika:

    • Nilai ini sama dengan atau melebihi spark.sql.adaptive.advisory PartitionSizeInBytes

    • Semua ukuran partisi berada dalam batas ini

    Mengganti pengaturan spark.sql.join.preferSortMergeJoin .

    0 byte

    spark.sql.adaptive.optimalkan SkewsInRebalancePartitions.enabled

    Menentukan apakah akan mengoptimalkan partisi shuffle miring dengan membaginya menjadi partisi yang lebih kecil berdasarkan spark.sql.adaptive.advisory. PartitionSizeInBytes Membutuhkan spark.sql.adaptive.enabled agar benar.

    true

    spark.sql.adaptive.rebalance PartitionsSmallPartitionFactor

    Mendefinisikan faktor ambang ukuran untuk menggabungkan partisi selama pemisahan. Partisi yang lebih kecil dari faktor ini dikalikan dengan spark.sql. PartitionSizeInBytes adaptive.advisory digabungkan.

    0.2

    spark.sql.adaptive.miring Join.enabled

    Menentukan apakah akan menangani kemiringan data dalam gabungan yang diacak dengan memisahkan dan secara opsional mereplikasi partisi miring. Berlaku untuk gabungan hash sort-merge dan shuffled. Membutuhkan spark.sql.adaptive.enabled agar benar.

    true

    spark.sql.adaptive.miring Join.skewedPartitionFactor

    Menentukan faktor ukuran yang menentukan kemiringan partisi. Partisi miring ketika ukurannya melebihi keduanya:

    • Faktor ini dikalikan dengan ukuran partisi median

    • Nilai spark.sql.adaptive.skew Join.skewedPartitionThresholdInBytes

    5

    spark.sql.adaptive.miring Join.skewedPartitionThresholdInBytes

    Menetapkan ambang ukuran dalam byte untuk mengidentifikasi partisi miring. Partisi miring ketika ukurannya melebihi keduanya:

    • Ambang batas ini

    • Ukuran partisi median dikalikan dengan spark.sql.adaptive.skew Join.skewedPartitionFactor

    Sebaiknya setel nilai ini lebih besar dari spark. PartitionSizeInBytes sql.adaptive.advisory.

    256MB

    Spark.sql.BroadcastTimeout

    Mengontrol periode batas waktu dalam detik untuk operasi siaran selama bergabung siaran.

    300 detik

    spark.sql.cbo.diaktifkan

    Menentukan apakah akan mengaktifkan optimasi berbasis biaya (CBO) untuk estimasi statistik rencana.

    false

    spark.sql.cbo.gabung Reorder.dp.star.filter

    Menentukan apakah akan menerapkan heuristik filter star-join selama enumerasi gabungan berbasis biaya.

    false

    spark.sql.cbo.gabung Reorder.dp.threshold

    Menetapkan jumlah maksimum node bergabung yang diizinkan dalam algoritma pemrograman dinamis.

    12

    spark.sql.cbo.gabung Reorder.enabled

    Menentukan apakah akan mengaktifkan pengaturan ulang gabungan dalam pengoptimalan berbasis biaya (CBO).

    false

    spark.sql.cbo.plan Stats.enabled

    Menentukan apakah akan mengambil jumlah baris dan statistik kolom dari katalog selama pembuatan rencana logis.

    false

    spark.sql.cbo.star SchemaDetection

    Menentukan apakah akan mengaktifkan penataan ulang gabungan berdasarkan deteksi skema bintang.

    false

    spark.sql.files.max PartitionNum

    Menetapkan target jumlah maksimum partisi file terpisah untuk sumber berbasis file (Parquet, JSON, dan ORC). Mengubah skala partisi ketika jumlah awal melebihi nilai ini. Ini adalah target yang disarankan, bukan batas yang dijamin.

    (tidak ada)

    spark.sql.files.max RecordsPerFile

    Menetapkan jumlah maksimum catatan untuk menulis ke satu file. Tidak ada batasan yang berlaku ketika disetel ke nol atau nilai negatif.

    0

    spark.sql.files.min PartitionNum

    Menetapkan target jumlah minimum partisi file terpisah untuk sumber berbasis file (Parquet, JSON, dan ORC). Defaultnya spar NodeDefaultParallelism k.sql.leaf. Ini adalah target yang disarankan, bukan batas yang dijamin.

    (tidak ada)

    spark.sql.in MemoryColumnarStorage.batchSize

    Mengontrol ukuran batch untuk cache kolumnar. Meningkatkan ukuran meningkatkan pemanfaatan memori dan kompresi tetapi meningkatkan risiko kesalahan kehabisan memori.

    10000

    spark.sql.in MemoryColumnarStorage.compressed

    Menentukan apakah akan secara otomatis memilih codec kompresi untuk kolom berdasarkan statistik data.

    true

    spark.sql.in MemoryColumnarStorage.enableVectorizedReader

    Menentukan apakah akan mengaktifkan pembacaan vektor untuk cache kolumnar.

    true

    spark.sql.legacy.allow HashOnMapType

    Menentukan apakah akan mengizinkan operasi hash pada struktur data tipe peta. Pengaturan lama ini mempertahankan kompatibilitas dengan penanganan jenis peta versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.allow NegativeScaleOfDecimal

    Menentukan apakah akan mengizinkan nilai skala negatif dalam definisi tipe desimal. Pengaturan lama ini mempertahankan kompatibilitas dengan versi Spark lama yang mendukung skala desimal negatif.

    (tidak ada)

    spark.sql.legacy.cast ComplexTypesToString.enabled

    Menentukan apakah akan mengaktifkan perilaku lama untuk mentransmisikan tipe kompleks ke string. Mempertahankan kompatibilitas dengan aturan konversi tipe versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.char VarcharAsString

    Menentukan apakah akan memperlakukan tipe CHAR dan VARCHAR sebagai tipe STRING. Pengaturan lama ini menyediakan kompatibilitas dengan penanganan jenis string versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.create EmptyCollectionUsingStringType

    Menentukan apakah akan membuat koleksi kosong menggunakan elemen tipe string. Pengaturan lama ini mempertahankan kompatibilitas dengan perilaku inisialisasi koleksi versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.eksponent LiteralAsDecimal.enabled

    Menentukan apakah akan menafsirkan literal eksponensial sebagai tipe desimal. Pengaturan lama ini mempertahankan kompatibilitas dengan penanganan literal numerik versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.json.allow EmptyString.enabled

    Menentukan apakah akan mengizinkan string kosong dalam pemrosesan JSON. Pengaturan lama ini mempertahankan kompatibilitas dengan perilaku penguraian JSON versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.parquet.int96 RebaseModeInRead

    Menentukan apakah akan menggunakan mode rebase stempel waktu INT96 lama saat membaca file Parquet. Pengaturan lama ini mempertahankan kompatibilitas dengan penanganan stempel waktu versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.legacy.waktu ParserPolicy

    Mengontrol perilaku penguraian waktu untuk kompatibilitas mundur. Pengaturan lama ini menentukan bagaimana cap waktu dan tanggal diurai dari string.

    (tidak ada)

    spark.sql.legacy.type Coercion.datetimeToString.enabled

    Menentukan apakah akan mengaktifkan perilaku paksaan tipe lama saat mengonversi nilai datetime menjadi string. Mempertahankan kompatibilitas dengan aturan konversi datetime versi Spark yang lebih lama.

    (tidak ada)

    spark.sql.maks SinglePartitionBytes

    Menetapkan ukuran partisi maksimum dalam byte. Perencana memperkenalkan operasi shuffle untuk partisi yang lebih besar untuk meningkatkan paralelisme.

    128m

    spark.sql.metadataCacheTtlSeconds

    Mengontrol waktu hidup (TTL) untuk cache metadata. Berlaku untuk metadata file partisi dan cache katalog sesi. Membutuhkan:

    • Nilai positif lebih besar dari nol

    • spark.sql.catalogImplementasi disetel ke sarang

    • spark.sql. PartitionFileCacheSize hive.filesource lebih besar dari nol

    • spark.sql. FilesourcePartitions hive.manage disetel ke true

    -1000ms

    spark.sql.optimizer.kolaps ProjectAlwaysInline

    Menentukan apakah akan menciutkan proyeksi yang berdekatan dan ekspresi sebaris, bahkan ketika itu menyebabkan duplikasi.

    false

    spark.sql.optimizer.dynamic PartitionPruning.enabled

    Menentukan apakah akan menghasilkan predikat untuk kolom partisi yang digunakan sebagai kunci gabungan.

    true

    spark.sql.optimizer.enable CsvExpressionOptimization

    Menentukan apakah akan mengoptimalkan ekspresi CSV di pengoptimal SQL dengan memangkas kolom yang tidak perlu dari operasi from_csv.

    true

    spark.sql.optimizer.enable JsonExpressionOptimization

    Menentukan apakah akan mengoptimalkan ekspresi JSON di pengoptimal SQL dengan:

    • Memangkas kolom yang tidak perlu dari operasi from_json

    • Menyederhanakan kombinasi from_json dan to_json

    • Mengoptimalkan operasi named_struct

    true

    spark.sql.optimizer.excludedRules

    Mendefinisikan aturan pengoptimal untuk dinonaktifkan, diidentifikasi oleh nama aturan yang dipisahkan koma. Beberapa aturan tidak dapat dinonaktifkan karena diperlukan untuk kebenaran. Pengoptimal mencatat aturan mana yang berhasil dinonaktifkan.

    (tidak ada)

    spark.sql.optimizer.runtime.bloom Filter.applicationSideScanSizeThreshold

    Menetapkan ukuran pemindaian agregat minimum dalam byte yang diperlukan untuk menyuntikkan filter Bloom di sisi aplikasi.

    10GB

    spark.sql.optimizer.runtime.bloom Filter.creationSideThreshold

    Mendefinisikan ambang batas ukuran maksimum untuk menyuntikkan filter Bloom di sisi pembuatan.

    10MB

    spark.sql.optimizer.runtime.bloom Filter.enabled

    Menentukan apakah akan menyisipkan filter Bloom untuk mengurangi data shuffle ketika salah satu sisi penggabungan shuffle memiliki predikat selektif.

    true

    spark.sql.optimizer.runtime.bloom Filter.expectedNumItems

    Mendefinisikan jumlah default item yang diharapkan dalam filter Bloom runtime.

    1000000

    spark.sql.optimizer.runtime.bloom Filter.maxNumBits

    Menetapkan jumlah maksimum bit yang diizinkan dalam filter Bloom runtime.

    67108864

    spark.sql.optimizer.runtime.bloom Filter.maxNumItems

    Menetapkan jumlah maksimum item yang diharapkan yang diizinkan dalam filter Bloom runtime.

    4000000

    spark.sql.optimizer.runtime.bloom Filter.numBits

    Mendefinisikan jumlah bit default yang digunakan dalam runtime filter Bloom.

    8388608

    spark.sql.optimizer.runtime.row LevelOperationGroupFilter.enabled

    Menentukan apakah akan mengaktifkan pemfilteran grup runtime untuk operasi tingkat baris. Memungkinkan sumber data untuk:

    • Pangkas seluruh kelompok data (seperti file atau partisi) menggunakan filter sumber data

    • Jalankan kueri runtime untuk mengidentifikasi catatan yang cocok

    • Buang grup yang tidak perlu untuk menghindari penulisan ulang yang mahal

    Pembatasan:

    • Tidak semua ekspresi dapat dikonversi ke filter sumber data

    • Beberapa ekspresi memerlukan evaluasi Spark (seperti subquery)

    true

    spark.sql.optimizer.runtime Filter.number.threshold

    Menetapkan jumlah total filter runtime yang disuntikkan (non-DPP). Ini untuk mencegah OOM driver dengan terlalu banyak filter Bloom.

    10

    spark.sql.optimizer.runtime Filter.semiJoinReduction.enabled

    Menentukan apakah akan menyisipkan semi-join untuk mengurangi data shuffle ketika salah satu sisi shuffle join memiliki predikat selektif.

    false

    spark.sql.parquet.aggregatePushdown

    Menentukan apakah akan menekan agregat ke Parket untuk pengoptimalan. Mendukung:

    • MIN dan MAX untuk tipe boolean, integer, float, dan tanggal

    • HITUNG untuk semua tipe data

    Melempar pengecualian jika statistik hilang dari footer file Parquet mana pun.

    false

    spark.sql.parquet.columnar ReaderBatchSize

    Mengontrol jumlah baris di setiap batch pembaca vektor parket. Pilih nilai yang menyeimbangkan overhead kinerja dan penggunaan memori untuk mencegah kesalahan kehabisan memori.

    4096

    spark.sql.parquet.enable VectorizedReader

    Menentukan apakah akan mengaktifkan decoding parket vektor.

    true

    spark.sql.shuffle.partisi

    Menetapkan jumlah partisi default untuk pengocokan data selama gabungan atau agregasi. Tidak dapat dimodifikasi antara kueri streaming terstruktur yang dimulai ulang dari lokasi pos pemeriksaan yang sama.

    200

    spark.sql.dicocokkan HashJoinFactor

    Mendefinisikan faktor perkalian yang digunakan untuk menentukan kelayakan shuffle hash join. Gabungan hash shuffle dipilih ketika ukuran data sisi kecil dikalikan dengan faktor ini kurang dari ukuran data sisi besar.

    3

    spark.sql.sources.paralel PartitionDiscovery.threshold

    Menetapkan jumlah maksimum jalur untuk daftar file sisi driver dengan sumber berbasis file (Parquet, JSON, dan ORC). Ketika terlampaui selama penemuan partisi, file dicantumkan menggunakan pekerjaan terdistribusi Spark terpisah.

    32

    spark.sql.statistics.histogram.enabled

    Menentukan apakah akan menghasilkan histogram setinggi selama perhitungan statistik kolom untuk meningkatkan akurasi estimasi. Membutuhkan pemindaian tabel tambahan di luar yang diperlukan untuk statistik kolom dasar.

    false

    percikan dinamis Allocation.executorIdleTimeout

    Menetapkan durasi eksekutor harus diam sebelum dihapus saat alokasi dinamis diaktifkan.

    60-an

    percikan dinamis Allocation.schedulerBacklogTimeout

    Menetapkan durasi tugas yang tertunda harus backlog sebelum pelaksana baru diminta saat alokasi dinamis diaktifkan.

    1 detik

    percikan dinamis Allocation.sustainedSchedulerBacklogTimeout

    Sama seperti spark.dynamicAllocation.schedulerBacklogTimeout, tetapi hanya digunakan untuk permintaan eksekutor berikutnya.

    (nilai spark.dynamicAllocation.schedulerBacklogTimeout)

    spark.scheduler.min RegisteredResourcesRatio

    Menetapkan rasio minimum sumber daya terdaftar (sumber daya terdaftar/total sumber daya yang diharapkan) untuk menunggu sebelum penjadwalan dimulai. Ditentukan sebagai ganda antara 0,0 dan 1,0. Terlepas dari apakah rasio minimum sumber daya telah tercapai, jumlah waktu maksimum yang akan menunggu sebelum penjadwalan dimulai dikendalikan oleh spark. RegisteredResourcesWaitingTime scheduler.max.

    0.8

    spark.scheduler.max RegisteredResourcesWaitingTime

    Menetapkan jumlah waktu maksimum untuk menunggu sumber daya mendaftar sebelum penjadwalan dimulai.

    30-an

    spark.sql.hive.metastore PartitionPruningFallbackOnException

    Menentukan apakah akan kembali mendapatkan semua partisi dari Hive metastore dan melakukan pemangkasan partisi di sisi klien Spark saat bertemu MetaException dari metastore.

    false

    spark.sql.cross Join.enabled

    Menentukan apakah akan mengizinkan kueri yang berisi produk cartesian tanpa sintaks CROSS JOIN eksplisit.

    true

    spark.sql.analyzer.maxIterations

    Menetapkan jumlah maksimum iterasi yang dijalankan oleh query Analyzer sebelum menyerah. Nilai yang lebih tinggi memungkinkan penganalisis untuk memproses kueri yang sangat besar atau sangat bersarang.

    100

    spark.sql.dataprefetch.filescan.max ParallelismPerTask

    Menetapkan jumlah maksimum pemisahan file untuk pra-pengambilan secara bersamaan untuk setiap tugas saat memindai file.

    4

    spark.sql.iceberg.data-prefetch.enabled

    Menentukan apakah akan mengaktifkan pengoptimalan pra-pengambilan data saat membaca Iceberg tabel.

    true

    spark.sql.legacy.null ValueWrittenAsQuotedEmptyStringCsv

    Menentukan apakah akan mengembalikan perilaku lama penulisan nol sebagai string kosong yang dikutip dalam keluaran CSV. Jika salah, Spark menulis nol sebagai string kosong yang tidak dikutip.

    false

    percikan maks RemoteBlockSizeFetchToMem

    Menetapkan ambang ukuran di mana Spark mengambil blok jarak jauh ke disk alih-alih memori. Ini menghindari satu permintaan besar yang menghabiskan terlalu banyak memori.

    200m

    spark.emr-serverless.allocation.batch.size

    Menetapkan jumlah pelaksana untuk diminta sekaligus di setiap putaran alokasi pelaksana.

    20

    Nama Properti Deskripsi nilai default

    spark.sql.auto BroadcastJoinThreshold

    Menetapkan ukuran tabel maksimum dalam byte untuk penyiaran ke node pekerja selama bergabung. Setel ke -1 untuk menonaktifkan penyiaran.

    10MB

    spark.io.kompressi.codec

    Menetapkan codec yang digunakan untuk mengompres data internal seperti partisi RDD, log peristiwa, variabel siaran, dan keluaran shuffle. Nilai yang didukung: lz4, snappy, zstd, gzip.

    lz4

    Spark.sql.session.Zona waktu

    Mendefinisikan zona waktu sesi untuk menangani stempel waktu dalam literal string dan konversi objek Java. Menerima:

    • Region-based ID dalam area/city format (seperti America/Los _Angeles)

    • Zona offset dalam HH:mm:ss format (+/-) HH, (+/-)HH:mm, atau (+/-) (seperti -08 atau + 01:00)

    • UTC atau Z sebagai alias untuk + 00:00

    (nilai zona waktu lokal)

    spark.cleanrooms.executor.memory OverheadFactor

    Menetapkan fraksi dari total memori eksekutor yang digunakan untuk menentukan pemisahan antara spark.executor.memory dan spark.executor.memoryOverhead. Ditentukan sebagai ganda antara 0,0 dan kurang dari 1,0.

    0.1

    spark.cleanrooms.driver.memory OverheadFactor

    Menetapkan fraksi total memori driver yang digunakan untuk menentukan pemisahan antara spark.driver.memory dan spark.driver.memoryOverhead. Ditentukan sebagai ganda antara 0,0 dan kurang dari 1,0.

    0.1

    Spark.Memory.StorageFraction

    Menetapkan jumlah memori penyimpanan yang kebal terhadap penggusuran, dinyatakan sebagai sebagian kecil dari ukuran wilayah yang disisihkan oleh spark.memory.fraction. Semakin tinggi ini, semakin sedikit memori kerja yang tersedia untuk eksekusi dan tugas dapat tumpah ke disk lebih sering. Dianjurkan untuk meninggalkan ini pada nilai default.

    0.5

    Spark.rpc.AskTimeout

    Menetapkan durasi operasi permintaan RPC untuk menunggu sebelum waktu habis.

    (nilai spark.network.timeout)

    Spark.Executor.HeartbeatInterval

    Mengatur interval antara detak jantung setiap pelaksana kepada pengemudi. Heartbeats memberi tahu pengemudi bahwa eksekutor masih hidup dan memperbaruinya dengan metrik untuk tugas yang sedang berjalan. spark.executor.heartbeatInterval harus jauh lebih kecil dari spark.network.timeout.

    10s

    spark.stage.max ConsecutiveAttempts

    Menetapkan jumlah upaya tahap berturut-turut yang diizinkan sebelum tahap dibatalkan.

    4

    spark.task.cpus

    Menetapkan jumlah core yang akan dialokasikan untuk setiap tugas.

    1

    spark.shuffle.file.buffer

    Menetapkan ukuran buffer dalam memori untuk setiap aliran keluaran file shuffle, dalam KiB kecuali ditentukan lain. Buffer ini mengurangi jumlah pencarian disk dan panggilan sistem yang dilakukan dalam membuat file shuffle perantara.

    32k

    spark.reducer.max SizeInFlight

    Menetapkan ukuran maksimum output peta untuk diambil secara bersamaan dari setiap tugas pengurangan, di MiB kecuali ditentukan lain. Karena setiap output memerlukan buffer untuk menerimanya, ini mewakili overhead memori tetap per tugas mengurangi, jadi jaga agar tetap kecil kecuali Anda memiliki sejumlah besar memori.

    48m

  9. (Opsional) Untuk Compute payer, pilih anggota kolaborasi yang membayar biaya komputasi pekerjaan.

    catatan

    Jika hanya ada satu kandidat pembayar untuk perhitungan pekerjaan dalam kolaborasi, itu default ke pembayar itu.

  10. Pilih Jalankan.

    catatan

    Anda tidak dapat menjalankan pekerjaan jika anggota yang dapat menerima hasil belum mengonfigurasi pengaturan hasil pekerjaan.

  11. Terus sesuaikan parameter dan jalankan pekerjaan Anda lagi, atau pilih tombol + untuk memulai pekerjaan baru di tab baru.