PyODPS DataFrame memungkinkan Anda memperluas komputasi bawaan dengan fungsi yang ditentukan pengguna (user-defined functions/UDF) dan paket Python pihak ketiga. Topik ini mencakup pemetaan elemen per elemen menggunakan map, transformasi tingkat baris dan agregasi kustom menggunakan apply, referensi resource di dalam UDF, serta cara mengunggah dan mengonfigurasi paket pihak ketiga.
Prasyarat
Sebelum memulai, pastikan Anda telah:
-
Objek DataFrame yang dibuat dari tabel MaxCompute atau pandas DataFrame
-
Mengaktifkan dukungan Python UDF di proyek MaxCompute Anda (diperlukan untuk
mapdanapplydengan fungsi Python)
Layanan publik Alibaba Cloud tidak mendukung Python UDF. Jika proyek Anda tidak mendukung Python UDF, metode map dan fungsi bawaan yang bergantung padanya tidak tersedia.
Batasan yang diketahui
| Batasan | Detail |
|---|---|
| Tipe yang tidak didukung | Metode map dan apply tidak menerima tipe LIST atau DICT sebagai input maupun output. |
| Pustaka biner pra-instal | Satu-satunya pustaka pihak ketiga pra-instal yang berisi kode C adalah NumPy. Semua pustaka biner lainnya memerlukan pengunggahan eksplisit. |
| Kompatibilitas Python 2/3 | Karena perbedaan byte code antar versi Python, kode yang menggunakan sintaks khusus Python 3 (seperti yield from) dapat gagal pada Worker MaxCompute yang menjalankan Python 2.7. Pastikan kode Anda berjalan dengan benar sebelum diterapkan ke lingkungan produksi menggunakan API MapReduce di Python 3. |
| Platform pembuatan paket biner | File wheel yang dibuat di macOS atau Windows tidak dapat digunakan di MaxCompute. Buat paket biner di shell Linux. |
Terapkan UDF ke kolom
Gunakan metode map pada objek Sequence untuk memanggil UDF pada setiap elemen.
>>> iris.sepallength.map(lambda x: x + 1).head(5)
sepallength
0 6.1
1 5.9
2 5.7
3 5.6
4 6.0
Jika tipe Sequence berubah setelah map, tentukan secara eksplisit tipe baru tersebut:
>>> iris.sepallength.map(lambda x: 't' + str(x), 'string').head(5)
sepallength
0 t5.1
1 t4.9
2 t4.7
3 t4.6
4 t5.0
Hindari bug penangkapan variabel closure
Ketika UDF berisi closure, perubahan eksternal pada variabel yang ditangkap akan memengaruhi perilaku fungsi. Kode berikut menghasilkan hasil yang tidak diinginkan — setiap SequenceExpr dalam dfs berakhir menjadi df.sepal_length + 9:
>>> dfs = []
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(lambda x: x + i))
Perbaiki hal ini dengan mengembalikan lambda dari fungsi luar, atau dengan menggunakan functools.partial:
# Opsi 1: gunakan fungsi factory
>>> dfs = []
>>> def get_mapper(i):
>>> return lambda x: x + i
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(get_mapper(i)))# Opsi 2: gunakan functools.partial
>>> import functools
>>> dfs = []
>>> for i in range(10):
>>> dfs.append(df.sepal_length.map(functools.partial(lambda v, x: x + v, i)))
Gunakan UDF yang sudah ada
Berikan nama fungsi (string) atau objek Function ke map untuk memanggil UDF yang sudah ada. Untuk detailnya, lihat Functions.
Lacak eksekusi dengan pencacah
Gunakan get_execution_context untuk mengakses pencacah dari dalam UDF. Nilai pencacah muncul di JSONSummary LogView.
from odps.udf import get_execution_context
def h(x):
ctx = get_execution_context()
counters = ctx.get_counters()
counters.get_counter('df', 'add_one').increment(1)
return x + 1
df.field.map(h)
Terapkan UDF ke baris
Gunakan apply dengan axis=1 untuk memanggil UDF pada setiap baris. UDF menerima satu baris dalam satu waktu; ambil nilai bidang berdasarkan nama atribut atau indeks.
Kembalikan satu nilai per baris
Atur reduce=True untuk mengembalikan Sequence. Tentukan tipe output dengan parameter types (default adalah STRING).
>>> iris.apply(lambda row: row.sepallength + row.sepalwidth, axis=1, reduce=True, types='float').rename('sepaladd').head(3)
sepaladd
0 8.6
1 7.9
2 7.9
Kembalikan beberapa baris menggunakan yield
Atur reduce=False dan gunakan yield untuk menghasilkan beberapa baris per baris input. Tentukan nama dan tipe bidang output dengan names dan types.
>>> iris.count()
150
>>> def handle(row):
>>> yield row.sepallength - row.sepalwidth, row.sepallength + row.sepalwidth
>>> yield row.petallength - row.petalwidth, row.petallength + row.petalwidth
>>> iris.apply(handle, axis=1, names=['iris_add', 'iris_sub'], types=['float', 'float']).count()
300
Anotasikan skema output langsung pada fungsi untuk menghindari pengulangan saat pemanggilan:
>>> from odps.df import output
>>> @output(['iris_add', 'iris_sub'], ['float', 'float'])
>>> def handle(row):
>>> yield row.sepallength - row.sepalwidth, row.sepallength + row.sepalwidth
>>> yield row.petallength - row.petalwidth, row.petallength + row.petalwidth
>>> iris.apply(handle, axis=1).count()
300
Setara: map_reduce hanya map
map_reduce dalam mode hanya map setara dengan apply dengan axis=1:
>>> iris.map_reduce(mapper=handle).count()
300
Gunakan UDTF yang sudah ada
Untuk memanggil user-defined table-valued function (UDTF) yang sudah ada di MaxCompute, berikan namanya sebagai string:
>>> iris['name', 'sepallength'].apply('your_func', axis=1, names=['name2', 'sepallength2'], types=['string', 'float'])
Gabungkan output baris dengan lateral view
Ketika reduce=False, gabungkan output UDF dengan kolom asli menggunakan lateral view — berguna untuk agregasi:
>>> from odps.df import output
>>> @output(['iris_add', 'iris_sub'], ['float', 'float'])
>>> def handle(row):
>>> yield row.sepallength - row.sepalwidth, row.sepallength + row.sepalwidth
>>> yield row.petallength - row.petalwidth, row.petallength + row.petalwidth
>>> iris[iris.category, iris.apply(handle, axis=1)]
Terapkan agregasi kustom ke kolom
Gunakan apply dengan axis=0 (atau tanpa argumen axis) untuk meneruskan kelas agregasi kustom ke semua objek Sequence. Kelas tersebut harus mengimplementasikan buffer, __call__, merge, dan getvalue.
class Agg(object):
def buffer(self):
return [0.0, 0]
def __call__(self, buffer, val):
buffer[0] += val
buffer[1] += 1
def merge(self, buffer, pbuffer):
buffer[0] += pbuffer[0]
buffer[1] += pbuffer[1]
def getvalue(self, buffer):
if buffer[1] == 0:
return 0.0
return buffer[0] / buffer[1]>>> iris.exclude('name').apply(Agg)
sepallength_aggregation sepalwidth_aggregation petallength_aggregation petalwidth_aggregation
0 5.843333 3.054 3.758667 1.198667
Baca resource MaxCompute di dalam UDF
UDF dapat membaca resource MaxCompute — resource file dan resource tabel — atau mereferensikan objek Collection sebagai resource. Bungkus UDF dalam closure atau kelas callable agar resource dimuat sekali saat inisialisasi, bukan per baris.
Memuat resource di dalam closure (bukan pada setiap pemanggilan fungsi) menghindari overhead inisialisasi berulang — misalnya, saat memuat tabel lookup atau artefak model.
UDF tingkat baris dengan resource file dan collection
>>> file_resource = o.create_resource('pyodps_iris_file', 'file', file_obj='Iris-setosa')
>>> iris_names_collection = iris.distinct('name')[:2]
>>> iris_names_collection
sepallength
0 Iris-setosa
1 Iris-versicolor>>> def myfunc(resources): # resource diteruskan sesuai urutan pemanggilan
>>> names = set()
>>> fileobj = resources[0] # resource file direpresentasikan sebagai objek mirip file
>>> for l in fileobj:
>>> names.add(l)
>>> collection = resources[1]
>>> for r in collection:
>>> names.add(r.name) # ambil nilai berdasarkan nama bidang atau offset
>>> def h(x):
>>> if x in names:
>>> return True
>>> else:
>>> return False
>>> return h
>>> df = iris.distinct('name')
>>> df = df[df.name,
>>> df.name.map(myfunc, resources=[file_resource, iris_names_collection], rtype='boolean').rename('isin')]
>>> df
name isin
0 Iris-setosa True
1 Iris-versicolor True
2 Iris-virginica False
Saat membaca tabel partisi, bidang partisi tidak disertakan.
UDF tingkat baris dengan DataFrame lokal sebagai resource
Variabel lokal dapat direferensikan sebagai resource di MaxCompute saat eksekusi. Pada contoh berikut, stop_words adalah DataFrame lokal yang diteruskan oleh pelaksana ke UDF sebagai resource:
>>> words_df
sentence
0 Hello World
1 Hello Python
2 Life is short I use Python
>>> import pandas as pd
>>> stop_words = DataFrame(pd.DataFrame({'stops': ['is', 'a', 'I']}))
>>> @output(['sentence'], ['string'])
>>> def filter_stops(resources):
>>> stop_words = set([r[0] for r in resources[0]])
>>> def h(row):
>>> return ' '.join(w for w in row[0].split() if w not in stop_words),
>>> return h
>>> words_df.apply(filter_stops, axis=1, resources=[stop_words])
sentence
0 Hello World
1 Hello Python
2 Life short use Python
Untuk operasi baris (axis=1), gunakan closure fungsi atau kelas callable untuk memuat resource. Untuk agregasi kolom, gunakan metode__init__sebagai gantinya.
Unggah pustaka Python pihak ketiga
MaxCompute mendukung pengunggahan paket Python dalam format .whl, .egg, .zip, dan .tar.gz. Semua dependensi harus ditentukan secara eksplisit — mengabaikan dependensi menyebabkan error impor saat waktu proses.
Pilih jalur pengunggahan berdasarkan jenis paket:
| Jenis paket | Metode pengunggahan | Catatan |
|---|---|---|
| Pra-instal | Tidak diperlukan | Hanya NumPy |
| Python murni (tanpa kode terkompilasi, tanpa operasi file) | Unggah sebagai resource file .whl |
Berfungsi untuk paket seperti python-dateutil, pytz, six. Versi MaxCompute yang lebih baru juga mendukung paket dengan operasi file. |
| Biner (ekstensi C terkompilasi) | Unggah sebagai resource arsip .zip, aktifkan isolasi |
Memerlukan tag platform cp27-cp27m-manylinux1_x86_64; bangun di Linux |
Paket Python murni
Secara default, PyODPS mendukung pustaka pihak ketiga yang berisi kode Python murni tetapi tanpa operasi file. Contoh berikut mengunggah python-dateutil dan dependensinya six.
Langkah 1: Unduh paket dan dependensinya. Paket harus dibangun untuk Linux.
$ pip download python-dateutil -d /to/path/
Ini mengunduh six-1.10.0-py2.py3-none-any.whl dan python_dateutil-2.5.3-py2.py3-none-any.whl.
Langkah 2: Unggah kedua file sebagai resource menggunakan create_resource.
# Pastikan ekstensi nama file benar.
>>> odps.create_resource('six.whl', 'file', file_obj=open('six-1.10.0-py2.py3-none-any.whl', 'rb'))
>>> odps.create_resource('python_dateutil.whl', 'file', file_obj=open('python_dateutil-2.5.3-py2.py3-none-any.whl', 'rb'))
Langkah 3: Gunakan pustaka tersebut. Tentukan secara global melalui options.df.libraries, atau per eksekusi melalui parameter libraries.
# Konfigurasi global (berlaku untuk semua operasi DataFrame selanjutnya dalam sesi ini)
>>> from odps import options
>>> def get_year(t):
>>> from dateutil.parser import parse
>>> return parse(t).strftime('%Y')
>>> options.df.libraries = ['six.whl', 'python_dateutil.whl']
>>> df.datestr.map(get_year)
datestr
0 2016
1 2015# Konfigurasi per eksekusi (hanya berlaku untuk pemanggilan ini)
>>> def get_year(t):
>>> from dateutil.parser import parse
>>> return parse(t).strftime('%Y')
>>> df.datestr.map(get_year).execute(libraries=['six.whl', 'python_dateutil.whl'])
datestr
0 2016
1 2015
Paket biner (berisi kode terkompilasi)
Paket yang menyertakan ekstensi C terkompilasi (seperti SciPy atau pandas) memerlukan langkah tambahan:
-
File
.whlharus menggunakan tag platformcp27-cp27m-manylinux1_x86_64. -
Unggah file sebagai resource arsip, dengan ekstensi
.whldiubah menjadi.zip. -
Atur
odps.isolation.session.enablekeTrue, atau aktifkanisolationdi pengaturan proyek Anda.
# Unggah paket biner sebagai arsip dengan ekstensi .zip.
>>> odps.create_resource('scipy.zip', 'archive', file_obj=open('scipy-0.19.0-cp27-cp27m-manylinux1_x86_64.whl', 'rb'))
# Jika isolasi sudah diaktifkan di proyek Anda, opsi berikut bersifat opsional.
>>> options.sql.settings = {'odps.isolation.session.enable': True}
>>> def psi(value):
>>> # Impor pustaka pihak ketiga di dalam fungsi untuk menghindari error
>>> # yang disebabkan oleh perbedaan struktural antar sistem operasi.
>>> from scipy.special import psi
>>> return float(psi(value))
>>> df.float_col.map(psi).execute(libraries=['scipy.zip'])
Untuk membuat paket biner dari sumber, jalankan perintah berikut di shell Linux. File wheel yang dibuat di macOS atau Windows tidak kompatibel dengan MaxCompute.
python setup.py bdist_wheel
Unggah melalui Konsol MaxCompute
Sebagai alternatif API PyODPS, unggah paket menggunakan add archive di Konsol MaxCompute.
Langkah 1: Identifikasi file paket yang tepat untuk setiap dependensi.
Kebanyakan paket menyediakan file .whl untuk berbagai platform. Untuk paket biner, cari file yang memiliki cp27-cp27m-manylinux1_x86_64 dalam namanya. Untuk paket Python murni, file wheel py2.py3-none-any mana pun dapat digunakan.
Langkah 2: Verifikasi semua dependensi yang diperlukan. Tabel berikut mencantumkan dependensi untuk paket umum.
| Paket | Dependencies |
|---|---|
| pandas | NumPy, python-dateutil, pytz, six |
| SciPy | NumPy |
| scikit-learn | NumPy, SciPy |
NumPy sudah pra-instal. Unggah hanya python-dateutil, pytz, pandas, SciPy, scikit-learn, dan six.
Langkah 3: Unduh file paket. Tabel berikut mencantumkan file spesifik yang harus diunduh untuk setiap paket.
| Paket | Berkas untuk Diunduh | Upload Resource Name |
|---|---|---|
| python-dateutil | python-dateutil-2.6.0.zip | python-dateutil.zip |
| pytz | pytz-2017.2.zip | pytz.zip |
| six | six-1.11.0.tar.gz | six.tar.gz |
| pandas | pandas-0.20.2-cp27-cp27m-manylinux1_x86_64.zip | pandas.zip |
| SciPy | scipy-0.19.0-cp27-cp27m-manylinux1_x86_64.zip | scipy.zip |
| scikit-learn | scikit_learn-0.18.1-cp27-cp27m-manylinux1_x86_64.zip | sklearn.zip |
Langkah 4: Unggah setiap file. Untuk paket biner (pandas, SciPy, scikit-learn), ubah ekstensi .whl menjadi .zip sebelum mengunggah.
add archive python-dateutil.zip;
add archive pandas.zip;
Tentukan pustaka untuk eksekusi
Gunakan options.df.libraries untuk mengatur pustaka secara global untuk sesi, atau berikan parameter libraries langsung ke metode eksekusi untuk membatasi cakupannya pada satu pemanggilan.
# Global: berlaku untuk semua operasi DataFrame selanjutnya dalam sesi ini
>>> from odps import options
>>> options.df.libraries = ['six.whl', 'python_dateutil.whl']# Lokal: hanya berlaku untuk pemanggilan eksekusi ini
>>> df.apply(my_func, axis=1).to_pandas(libraries=['six.whl', 'python_dateutil.whl'])