Node PyODPS 3 memungkinkan Anda menulis pekerjaan MaxCompute dalam Python dan menjadwalkannya untuk dijalankan secara berkala di DataWorks.
Pendahuluan
PyODPS adalah SDK Python untuk MaxCompute yang memungkinkan Anda menulis pekerjaan, mengkueri tabel dan tampilan (views), serta mengelola resource MaxCompute dalam Python. Untuk informasi selengkapnya, lihat Ikhtisar PyODPS. Di DataWorks, node PyODPS memungkinkan Anda menjadwalkan dan menjalankan tugas Python bersama pekerjaan lainnya.
Catatan penggunaan
-
Untuk memanggil paket pihak ketiga dari node PyODPS pada kelompok sumber daya DataWorks, gunakan serverless resource group dengan custom image.
CatatanMetode ini tidak berlaku jika kode Anda mencakup user-defined function (UDF) yang mereferensikan paket pihak ketiga. Untuk prosedur yang benar, lihat Contoh UDF: Gunakan paket pihak ketiga dalam Python UDF.
-
Untuk melakukan upgrade versi PyODPS, Anda dapat menggunakan custom image dalam serverless resource group untuk menjalankan perintah
/home/tops/bin/pip3 install pyodps==0.12.1. Anda dapat mengganti0.12.1dengan versi PyODPS target. Untuk exclusive resource group for scheduling, gunakan O&M Assistant untuk menjalankan perintah yang sama. -
Jika pekerjaan PyODPS Anda perlu mengakses lingkungan jaringan tertentu, seperti sumber data atau layanan di VPC atau jaringan IDC, gunakan serverless resource group. Untuk informasi selengkapnya tentang cara menghubungkan serverless resource group ke lingkungan target, lihat Solusi konektivitas jaringan.
-
Untuk informasi selengkapnya tentang sintaksis PyODPS, lihat Dokumentasi PyODPS.
-
Node PyODPS tersedia dalam dua jenis: PyODPS 2 (Python 2) dan PyODPS 3 (Python 3). Buat jenis node yang sesuai dengan versi Python Anda.
-
Jika menjalankan pernyataan SQL dalam node PyODPS tidak menghasilkan alur data (data lineage) yang benar di Data Map, atur secara manual parameter penjadwalan DataWorks terkait dalam kode pekerjaan. Untuk informasi tentang cara melihat alur data, lihat Lihat alur data. Untuk informasi tentang cara mengatur parameter, lihat Atur petunjuk parameter waktu proses. Kode contoh berikut memperoleh parameter yang diperlukan saat waktu proses.
import os ... # get DataWorks scheduler runtime parameters skynet_hints = {} for k, v in os.environ.items(): if k.startswith('SKYNET_'): skynet_hints[k] = v ... # setting hints while submitting a job o.execute_sql('INSERT OVERWRITE TABLE XXXX SELECT * FROM YYYY WHERE ***', hints=skynet_hints) ... -
Log output untuk node PyODPS memiliki ukuran maksimum 4 MB. Hindari mencetak data dalam jumlah besar ke log. Fokuslah pada output informasi peringatan dan progres.
Batasan
-
Saat menjalankan node PyODPS pada exclusive resource group for scheduling, kami menyarankan agar data yang diproses secara lokal tidak melebihi 50 MB. Batasan ini ditetapkan berdasarkan spesifikasi exclusive resource group. Jika data lokal melebihi ambang batas sistem operasi, kesalahan OOM (Got Killed) dapat terjadi. Hindari menulis kode pemrosesan data berlebihan dalam node PyODPS.
-
Saat menjalankan node PyODPS pada serverless resource group, konfigurasikan jumlah CU yang sesuai berdasarkan volume data.
CatatanPada serverless resource group, satu pekerjaan mendukung maksimal
64 CUs. Namun, kami menyarankan tidak lebih dari16 CUsuntuk mencegah kekurangan resource yang dapat memengaruhi startup pekerjaan. -
Kesalahan Got killed menunjukkan bahwa proses melebihi batas memori. Hindari operasi data lokal. Pekerjaan SQL dan DataFrame (kecuali operasi to_pandas) yang dimulai melalui PyODPS tidak tunduk pada batasan ini.
-
Anda dapat menggunakan library NumPy dan pandas yang telah pra-instal dalam kode yang tidak melibatkan user-defined functions. Paket pihak ketiga lain yang berisi kode biner tidak didukung.
-
Karena alasan kompatibilitas,
options.tunnel.use_instance_tunneldiatur keFalsesecara default di DataWorks. Untuk mengaktifkaninstance tunnelsecara global, atur nilai ini secara manual menjadiTrue. -
Definisi bytecode berbeda antara versi minor Python 3, seperti Python 3.8 dan Python 3.7.
MaxCompute saat ini menggunakan Python 3.7. Jika Anda menggunakan sintaksis dari versi Python 3 lainnya, seperti
finally blockdari Python 3.8, kesalahan akan terjadi selama eksekusi. Kami menyarankan menggunakan Python 3.7. -
PyODPS 3 mendukung eksekusi pada serverless resource group. Untuk membeli dan menggunakannya, lihat Gunakan serverless resource group.
-
Anda tidak dapat mengonfigurasi beberapa pekerjaan Python untuk dijalankan secara konkuren dalam satu node PyODPS.
-
Untuk mencetak log dalam node PyODPS, gunakan
print. Penggunaanlogger.infotidak didukung.
Prasyarat
Kaitkan mesin komputasi MaxCompute dengan ruang kerja DataWorks Anda.
Prosedur
-
Kembangkan kode Anda pada halaman editor node PyODPS 3.
Contoh kode PyODPS 3
Setelah membuat node PyODPS, Anda dapat mengedit dan menjalankan kode Anda. Untuk informasi selengkapnya tentang sintaksis PyODPS, lihat Operasi dasar. Contoh berikut mencakup lima skenario umum. Pilih yang sesuai dengan kebutuhan Anda.
Titik masuk ODPS
Setiap node PyODPS di DataWorks menyertakan variabel titik masuk ODPS global,
odpsatauo, yang tidak perlu Anda definisikan secara manual.print(odps.exist_table('PyODPS_iris'))Eksekusi SQL
Anda dapat menjalankan pernyataan SQL dalam node PyODPS. Untuk informasi selengkapnya, lihat SQL.
-
Secara default,
instance tunneldinonaktifkan di DataWorks, sehinggainstance.open_readermenggunakan antarmuka Result dan mengembalikan maksimal 10.000 catatan. Anda dapat menggunakanreader.countuntuk mendapatkan jumlah catatan. Untuk mengiterasi seluruh data, nonaktifkanlimit. Pernyataan berikut mengaktifkaninstance tunnelsecara global dan menonaktifkanlimit.options.tunnel.use_instance_tunnel = True options.tunnel.limit_instance_tunnel = False # Disable the limit to read all data. with instance.open_reader() as reader: # All data can be read through Instance Tunnel. -
Anda juga dapat mengaktifkan
instance tunneluntuk satu panggilanopen_readerdengan menambahkantunnel=Trueke panggilanopen_reader. Anda juga dapat menambahkanlimit=Falseuntuk menonaktifkan pembatasanlimituntuk panggilan tersebut.# Use Instance Tunnel for this open_reader call and read all data. with instance.open_reader(tunnel=True, limit=False) as reader:
Atur parameter waktu proses
-
Anda dapat mengatur parameter waktu proses dengan menggunakan parameter
hints, yang merupakandict. Untuk informasi selengkapnya tentang hints, lihat Operasi SET.o.execute_sql('select * from PyODPS_iris', hints={'odps.sql.mapper.split.size': 16}) -
Jika Anda mengatur konfigurasi global menggunakan
sql.settings, parameter waktu proses terkait akan ditambahkan ke setiap eksekusi.from odps import options options.sql.settings = {'odps.sql.mapper.split.size': 16} o.execute_sql('select * from PyODPS_iris') # Hints are added based on the global settings.
Baca hasil eksekusi
Anda dapat memanggil
open_readerlangsung pada instans eksekusi SQL. Ini mendukung dua skenario:-
SQL mengembalikan data terstruktur.
with o.execute_sql('select * from dual').open_reader() as reader: for record in reader: # Process each record. -
Jika Anda mengeksekusi pernyataan SQL seperti
desc, Anda dapat menggunakan propertireader.rawuntuk memperoleh hasil mentah eksekusi SQL.with o.execute_sql('desc dual').open_reader() as reader: print(reader.raw)CatatanJika Anda menggunakan parameter penjadwalan kustom dan menjalankan node PyODPS 3 langsung dari UI, Anda harus menuliskan nilai waktu secara eksplisit karena node tidak dapat menggantikan variabel saat waktu proses.
DataFrame
Anda juga dapat memproses data menggunakan DataFrame (Deprecated).
-
Eksekusi
Di lingkungan DataWorks, operasi DataFrame memerlukan pemanggilan eksplisit ke metode eksekusi langsung.
from odps.df import DataFrame iris = DataFrame(o.get_table('pyodps_iris')) for record in iris[iris.sepal_width < 3].execute(): # Call an immediate execution method to process each record.Jika Anda perlu memicu eksekusi langsung saat mencetak, Anda harus mengaktifkan
options.interactive.from odps import options from odps.df import DataFrame options.interactive = True # Enable the option at the beginning. iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepal_width.sum()) # This triggers immediate execution. -
Cetak informasi detail
options.verbosediaktifkan secara default di DataWorks, sehingga informasi detail seperti URL Logview dicetak selama eksekusi.
Pengembangan kode PyODPS 3
Contoh berikut menunjukkan cara menggunakan node PyODPS:
-
Siapkan dataset dengan membuat tabel sampel pyodps_iris. Untuk informasi selengkapnya, lihat Proses data DataFrame.
-
Buat DataFrame. Untuk informasi selengkapnya, lihat Buat DataFrame dari tabel MaxCompute.
-
Masukkan kode berikut ke dalam node PyODPS dan jalankan.
from odps.df import DataFrame # Create a DataFrame from an ODPS table. iris = DataFrame(o.get_table('pyodps_iris')) print(iris.sepallength.head(5))
Jalankan pekerjaan PyODPS
-
Pada panel Run Configuration, di bagian Compute Resource, konfigurasikan Compute Resource, computing quota, dan DataWorks Resource Group.
Catatan-
Untuk mengakses sumber data melalui jaringan publik atau VPC, Anda harus menggunakan kelompok sumber daya penjadwalan yang telah lulus uji konektivitas dengan sumber data tersebut. Untuk informasi selengkapnya, lihat Solusi konektivitas jaringan.
-
Anda dapat mengonfigurasi informasi Image berdasarkan kebutuhan pekerjaan.
-
-
Pada kotak dialog parameter di bilah alat, pilih sumber data MaxCompute yang telah dibuat dan klik Run untuk menjalankan pekerjaan PyODPS.
-
-
Untuk menjalankan node secara berkala, konfigurasikan properti penjadwalannya sesuai kebutuhan bisnis Anda. Untuk informasi selengkapnya, lihat Konfigurasi penjadwalan node.
Berbeda dengan node SQL di DataWorks, node PyODPS tidak mengganti string seperti ${param_name} dalam kode. Sebagai gantinya, sebelum kode dieksekusi, sebuah dictionary bernama
argsditambahkan ke variabel global, dari mana Anda dapat mengambil parameter penjadwalan. Misalnya, jika Anda mengaturds=${yyyymmdd}di bagian Parameter, Anda dapat mengambil informasi parameter dalam kode sebagai berikut.print('ds=' + args['ds']) ds=20240930CatatanJika Anda perlu memperoleh partisi bernama
ds, Anda dapat menggunakan metode berikut.o.get_table('table_name').get_partition('ds=' + args['ds']) -
Setelah node dikonfigurasi, Anda harus menerapkannya. Untuk informasi selengkapnya, lihat Penerapan node dan alur kerja.
-
Setelah pekerjaan diterapkan, Anda dapat melihat statusnya di Operation Center. Untuk informasi selengkapnya, lihat Memulai Operation Center.
Jalankan node menggunakan peran terkait
Anda dapat mengaitkan peran RAM untuk menjalankan node, yang memungkinkan Anda menjalankan tugas node dengan peran RAM tertentu untuk kontrol izin detail halus dan manajemen keamanan.
Langkah berikutnya
FAQ tentang PyODPS: Masalah umum selama eksekusi PyODPS dan cara mengatasinya.