Tutorial dlt: Pipeline Ingestion Data Berbasis Python

# Membangun Pipeline EL Berbasis Python dengan dlt (data load tool) Sebagian besar tim data menghabiskan waktu yang tidak sedikit pada bagian pipeline yang kurang menarik: menarik data dari sebuah AP...

By Ruby Abdullah · · tutorial
dltData IngestionELTData PipelineData EngineeringPython

Membangun Pipeline EL Berbasis Python dengan dlt (data load tool)

Sebagian besar tim data menghabiskan waktu yang tidak sedikit pada bagian pipeline yang kurang menarik: menarik data dari sebuah API dan menempatkannya di warehouse dengan skema yang rapi. dlt (singkatan dari data load tool, dari dltHub) adalah pustaka Python ringan yang dibuat khusus untuk pekerjaan itu. Tutorial ini membahas cara kerjanya dan bagaimana ia melengkapi perkakas dbt serta Dagster yang mungkin sudah Anda gunakan.

Satu catatan soal penamaan, karena singkatannya rancu: dlt yang dibahas di sini adalah pustaka ingesti Python open-source dari dltHub. Ini bukan Delta Lake (format penyimpanan) dan bukan DLT milik PyTorch. Ketika kita menyebut dlt, yang dimaksud adalah pip install dlt.

Posisi dlt dalam Modern Data Stack

Modern analytics stack umumnya digambarkan sebagai EL + T: Extract dan Load data mentah dulu, lalu Transform di warehouse.

  • Extract + Load adalah ranah dlt. Ia membaca dari sumber (REST API, database, file, generator) dan menulis data ke warehouse tujuan, menangani inferensi skema, normalisasi, serta state inkremental untuk Anda.
  • Transform adalah tempat dbt bekerja. Setelah tabel mentah masuk ke warehouse, model dbt membentuknya menjadi mart yang bersih, teruji, dan terdokumentasi.
  • Orkestrasi adalah tempat Dagster (atau Airflow, cron, GitHub Actions) bekerja. Ia menentukan kapan extract-load dlt berjalan dan kapan transform dbt berjalan, serta menyambungkan dependensi di antaranya.

Jadi dlt tidak menggantikan dbt atau Dagster. Ia mengisi celah yang sengaja dibiarkan keduanya: memasukkan data mentah ke warehouse dengan andal, dalam Python murni, tanpa harus menyiapkan platform ingesti yang berat. Berbeda dari platform konektor seperti Fivetran atau Airbyte — bagus saat konektornya ada, canggung saat tidak ada — dlt membiarkan Anda menulis fungsi Python biasa yang menghasilkan (yield) record, mendekorasinya, lalu mendapat pipeline yang tangguh, sadar skema, dan inkremental. Jika Anda bisa memanggil sebuah API di Python, Anda bisa memuatnya dengan dlt.

Instalasi dan Penyiapan Proyek

dlt berjalan di mana saja Python 3.8+ berjalan. Pasang paket inti beserta extras untuk destinasi Anda.

# Pustaka inti

pip install dlt

Dengan dependensi destinasi tertentu

pip install "dlt[duckdb]" # analitik lokal, sangat baik untuk pengembangan

pip install "dlt[bigquery]" # Google BigQuery

pip install "dlt[postgres]" # PostgreSQL

pip install "dlt[snowflake]" # Snowflake

Disarankan bekerja di dalam virtual environment.

python -m venv .venv

source .venv/bin/activate # Windows: .venv\Scripts\activate

pip install "dlt[duckdb]"

Scaffolding dengan dlt init

dlt init menghasilkan pipeline awal yang sudah terhubung ke sumber dan destinasi. Bentuk umumnya adalah dlt init .
# Memulai pipeline yang memuat ke DuckDB

dlt init github duckdb

Perintah ini membuat struktur proyek seperti berikut:

.

├── .dlt/

│ ├── config.toml # konfigurasi non-rahasia (nama dataset, opsi)

│ └── secrets.toml # kredensial (di-gitignore)

├── githubpipeline.py # skrip pipeline Anda

└── requirements.txt

Berkas .dlt/secrets.toml otomatis ditambahkan ke .gitignore. Jangan pernah meng-commit-nya.

Konsep Inti: Resource, Source, Pipeline

dlt memiliki tiga blok bangunan yang perlu dipahami sebelum menulis kode.

  • Resource (@dlt.resource): fungsi yang menghasilkan data — biasanya baris atau halaman (page) record. Setiap resource dipetakan ke satu tabel di destinasi.
  • Source (@dlt.source): fungsi yang mengelompokkan beberapa resource terkait (misalnya semua endpoint dari satu API). Ia mengembalikan resource yang dikelolanya.
  • Pipeline (dlt.pipeline(...)): objek yang menghubungkan source ke destinasi, menjalankan pemuatan, dan melacak state (skema, cursor inkremental, riwayat run).

Pipeline Pertama: REST API Berpaginasi ke DuckDB

Mari memuat issue dari REST API berpaginasi ke DuckDB. Kita pakai REST API GitHub sebagai contoh berjalan karena paginasinya mewakili kebanyakan API. Perhatikan bagaimana resource, pengelompokan ala source, dan pipeline menyatu: dlt membuat berkas DuckDB, menyimpulkan skema dari data, lalu menulis barisnya tanpa DDL atau pembuatan tabel manual.

import dlt

from dlt.sources.helpers import requests

@dlt.resource(tablename="issues", writedisposition="replace")

def githubissues(repo: str = "dlt-hub/dlt"):

url = f"https://api.github.com/repos/{repo}/issues"

params = {"state": "all", "perpage": 100, "page": 1}

while True:

response = requests.get(url, params=params)

response.raiseforstatus()

page = response.json()

if not page:

break

yield page # yield seluruh halaman; dlt akan meratakannya

params["page"] += 1

pipeline = dlt.pipeline(

pipelinename="github",

destination="duckdb",

datasetname="githubdata",

)

loadinfo = pipeline.run(githubissues())

print(loadinfo)

Beberapa hal yang perlu dicatat:

  • dlt.sources.helpers.requests adalah pembungkus tipis di atas requests yang menambahkan retry dan timeout yang masuk akal. Anda juga boleh memakai requests biasa.
  • Kita yield satu halaman penuh (sebuah list). dlt otomatis mengiterasi list itu, sehingga tiap elemen menjadi satu baris.
  • writedisposition="replace" berarti tiap run membangun ulang tabel dari awal. Nanti kita perbaiki dengan pemuatan inkremental.

Memeriksa hasil

Setelah run, periksa apa yang terjadi melalui objek pipeline (pipeline.lasttrace untuk metrik run, pipeline.lasttrace.lastnormalizeinfo untuk detail per tahap) dan kueri data yang dimuat via SQL.

print(pipeline.lasttrace)

with pipeline.sqlclient() as client:

with client.executequery("SELECT count() AS n FROM issues") as cursor:

print(cursor.fetchall())

Untuk DuckDB Anda juga bisa membuka berkas .duckdb yang dihasilkan langsung di klien DuckDB apa pun dan menjalankan SQL seperti SELECT id, title, state FROM githubdata.issues ORDER BY createdat DESC LIMIT 10;.

Inferensi Skema Otomatis, Normalisasi, dan Evolusi

Di sinilah dlt membuktikan nilainya.

Inferensi

dlt membaca record Anda, menyimpulkan nama kolom dan tipenya, lalu membuat tabel destinasi untuk Anda. Key dict Python menjadi kolom; tipe disimpulkan dari nilai (dengan aturan koersi yang masuk akal untuk tanggal, desimal, dan sebagainya).

Normalisasi JSON bersarang

Respons API jarang berbentuk datar. Sebuah issue GitHub berisi objek user bersarang dan array labels. dlt menormalkannya secara otomatis:

  • Objek bersarang diratakan menjadi kolom berprefiks, mis. user.login menjadi userlogin.
  • List bersarang dipecah menjadi tabel anak yang terhubung kembali ke induk melalui foreign key yang dihasilkan.

Jadi payload issues dengan array labels menghasilkan tabel issues sekaligus tabel anak issueslabels, terhubung lewat key yang dihasilkan.

SELECT l.name, i.title

FROM issueslabels l

JOIN issues i ON l.dltparentid = i.dltid;

dlt menambahkan kolom pembukuan berprefiks dlt (dltid, dltloadid, dltparentid) agar relasi dan batch pemuatan ini dapat ditelusuri.

Evolusi

Saat sumber menambah field baru, dlt mendeteksinya pada run berikutnya dan otomatis menambahkan kolom ke tabel destinasi — tanpa mematahkan pemuatan. Field yang dihapus cukup menghasilkan NULL untuk baris baru. Evolusi skema inilah yang membuat dlt tahan terhadap perubahan API upstream yang jika tidak demikian akan mematahkan loader buatan tangan.

Anda bisa memeriksa skema yang disimpulkan kapan saja.

print(pipeline.defaultschema.toprettyyaml())

Write Disposition: replace, append, merge

write
disposition mengatur bagaimana data baru berinteraksi dengan yang sudah ada di tabel.
# replace: hapus dan muat ulang tabel tiap run (full refresh)

@dlt.resource(writedisposition="replace")

def dimcountry(): ...

append: tambah baris baru, tidak menyentuh yang lama (cocok untuk event/log yang imutabel)

@dlt.resource(writedisposition="append")

def pageviews(): ...

merge: upsert — perbarui baris yang ada berdasarkan key, sisipkan yang baru

@dlt.resource(writedisposition="merge", primarykey="id")

def orders(): ...

  • replace paling sederhana dan tepat untuk tabel dimensi kecil yang masih layak dibangun ulang.
  • append untuk fakta yang tidak pernah berubah, seperti log event. Hati-hati duplikat saat run ulang.
  • merge adalah upsert. Anda deklarasikan primarykey (atau mergekey); dlt menghapus baris yang cocok dan menyisipkan batch masuk, sehingga run ulang bersifat idempoten.

Gunakan mergekey alih-alih primarykey ketika natural key terdiri dari beberapa kolom atau ketika Anda ingin mengatur predikat penghapusan secara terpisah dari identitas baris.

Pemuatan Inkremental dengan dlt.sources.incremental

Muat ulang penuh tidak skalabel. Pemuatan inkremental hanya mengambil record yang berubah sejak run terakhir, menggunakan cursor field (kolom yang naik secara monoton seperti updatedat atau id).

import dlt

from dlt.sources.helpers import requests

@dlt.resource(tablename="issues", writedisposition="merge", primarykey="id")

def githubissues(

repo: str = "dlt-hub/dlt",

updatedat=dlt.sources.incremental(

"updatedat",

initialvalue="2024-01-01T00:00:00Z",

),

):

url = f"https://api.github.com/repos/{repo}/issues"

params = {

"state": "all",

"perpage": 100,

"page": 1,

"since": updatedat.lastvalue, # minta hanya data baru ke API

"sort": "updated",

"direction": "asc",

}

while True:

response = requests.get(url, params=params)

response.raiseforstatus()

page = response.json()

if not page:

break

yield page

params["page"] += 1

Perilakunya:

  • Pada run pertama, updatedat.lastvalue sama dengan initialvalue. dlt memuat semuanya sejak titik itu.
  • dlt menyimpan nilai updatedat terbesar yang dilihatnya di state pipeline.
  • Pada run berikutnya, lastvalue adalah nilai maksimum tersimpan itu, sehingga Anda hanya mengambil record yang lebih baru.
  • Dikombinasikan dengan writedisposition="merge" dan primarykey="id", record yang diperbarui di upstream akan di-upsert, bukan diduplikasi.

dlt juga mendeduplikasi baris pada nilai batas (record dengan lastvalue persis) agar tepi data tidak termuat ganda. State dipertahankan di destinasi, sehingga bertahan lintas proses run yang terpisah — tanpa perlu penyimpanan state eksternal.

Ekstraksi Berbasis Konfigurasi dengan REST API Source

Menulis loop paginasi secara manual itu berulang dan melelahkan. dlt menyediakan helper deklaratif restapisource yang mendeskripsikan API sebagai konfigurasi. Ia menangani paginasi, autentikasi, dan penyambungan inkremental untuk Anda.

import dlt

from dlt.sources.restapi import restapisource

source = restapisource({

"client": {

"baseurl": "https://api.github.com/",

"auth": {

"type": "bearer",

"token": dlt.secrets["githubtoken"],

},

"paginator": {

"type": "headerlink", # GitHub memakai Link header

},

},

"resourcedefaults": {

"primarykey": "id",

"writedisposition": "merge",

},

"resources": [

{

"name": "issues",

"endpoint": {

"path": "repos/dlt-hub/dlt/issues",

"params": {

"state": "all",

"since": {

"type": "incremental",

"cursorpath": "updatedat",

"initialvalue": "2024-01-01T00:00:00Z",

},

},

},

},

{

"name": "comments",

"endpoint": {"path": "repos/dlt-hub/dlt/issues/comments"},

},

],

})

pipeline = dlt.pipeline(

pipelinename="githubrest",

destination="duckdb",

datasetname="githubdata",

)

print(pipeline.run(source))

Bentuk deklaratif adalah titik awal yang disarankan untuk REST API pada umumnya; turun ke @dlt.resource tulisan tangan hanya ketika sebuah API melakukan sesuatu yang tidak biasa.

Secrets dan Konfigurasi

dlt menjaga konfigurasi dan secrets tetap di luar kode Anda. Ia menyelesaikan nilai dari .dlt/secrets.toml, .dlt/config.toml, dan environment variable, dengan urutan spesifisitas tersebut.

.dlt/secrets.toml (di-gitignore):
[sources.github]

githubtoken = "ghppersonalaccesstokenanda"

[destination.bigquery]

location = "EU"

[destination.bigquery.credentials]

projectid = "my-gcp-project"

privatekey = "-----BEGIN PRIVATE KEY-----\n...\n-----END PRIVATE KEY-----\n"

clientemail = "loader@my-gcp-project.iam.gserviceaccount.com"

.dlt/config.toml (aman di-commit — tanpa secrets):
[runtime]

loglevel = "INFO"

[load]

workers = 4

Di produksi, utamakan environment variable. dlt memetakan hierarki TOML ke nama huruf besar yang dipisahkan garis bawah ganda.

export SOURCESGITHUBGITHUBTOKEN="ghp..."

export DESTINATIONBIGQUERYCREDENTIALS_PROJECTID="my-gcp-project"

Rujuk sebuah secret di kode dengan dlt.secrets["..."] atau, yang lebih idiomatik, dengan memberi fungsi Anda argumen bertipe dengan nilai default dlt.secrets.value — dlt menyuntikkan nilai yang sudah diselesaikan secara otomatis saat pemanggilan.

@dlt.source

def githubsource(githubtoken: str = dlt.secrets.value):

...

Data Contract dan Kontrol Skema

Secara default dlt mengevolusi skema dengan bebas, yang nyaman saat pengembangan tetapi berisiko di produksi. Pengaturan schemacontract memungkinkan Anda membatasi perilaku itu per tabel, kolom, dan tipe data.

@dlt.resource(

schemacontract={

"tables": "evolve", # izinkan tabel baru

"columns": "freeze", # tolak kolom baru -> error pada field tak terduga

"datatype": "freeze" # tolak perubahan tipe

}

)

def orders():

...

Mode yang tersedia adalah evolve (izinkan, default), freeze (lempar error), discardrow (buang baris bermasalah), dan discardvalue (buang nilai bermasalah tetapi pertahankan baris). Pola produksi yang umum adalah evolve untuk kolom tetapi freeze untuk tipe data, sehingga perubahan aditif lolos sementara pergeseran tipe terdeteksi sejak dini.

Transformer dan Paralelisme

@dlt.transformer mengonsumsi keluaran sebuah resource dan menghasilkan resource turunan — berguna untuk fan-out, seperti mengambil detail tiap ID yang dikembalikan endpoint daftar.
@dlt.resource

def issueids():

yield from [{"id": 1}, {"id": 2}, {"id": 3}]

@dlt.transformer(datafrom=issueids)

def issuedetail(issue):

resp = requests.get(f"https://api.example.com/issues/{issue['id']}")

yield resp.json()

Salurkan resource ke transformer

pipeline.run(issueids | issuedetail)

Untuk throughput, tandai resource parallelized=True dan setel jumlah worker per tahap (extract, normalize, load) di .dlt/config.toml.

@dlt.resource(parallelized=True)

def heavyendpoint():

...

Mulailah dengan nilai default dan naikkan jumlah worker hanya ketika sebuah tahap terukur menjadi bottleneck.

Menjalankan dan Men-deploy

Pipeline dlt hanyalah skrip Python, sehingga deployment-nya adalah apa pun yang menjalankan Python secara terjadwal.

  • Cron — opsi paling sederhana: 0 .venv/bin/python githubpipeline.py.
  • GitHub Actions — dlt deploy githubpipeline.py github-action --schedule "0 *" men-scaffold workflow dan menampilkan secret repositori yang perlu disetel (ia tidak pernah menyematkan kredensial).
  • Airflow / Dagster — panggil pipeline dari dalam sebuah task atau asset. Pola Dagster yang umum: sebuah asset menjalankan extract-load dlt, asset dbt di hilir menransformasinya, dan Dagster mengelola dependensi serta penjadwalan.

# Sketsa asset Dagster yang membungkus dlt

from dagster import asset

@asset

def rawgithubissues():

pipeline = dlt.pipeline(

pipelinename="github", destination="bigquery", datasetname="raw",

)

return pipeline.run(githubsource())

Inilah persis sambungan tempat dlt menyerahkan ke dbt: dlt mendaratkan raw, dbt membangun mart di atasnya.

Praktik Terbaik

  • Kembangkan dengan DuckDB, deploy ke warehouse Anda. DuckDB tidak butuh infrastruktur dan berjalan lokal, jadi beriterasilah cepat di sana, lalu ubah hanya argumen destination untuk produksi.
  • Utamakan merge + inkremental untuk apa pun yang besar. replace penuh layak untuk dimensi kecil tetapi boros dan lambat untuk tabel besar yang terus tumbuh.
  • Pilih cursor field yang tepat. Gunakan kolom yang naik secara andal (updatedat, id auto-increment). Hindari field yang bisa diisi ulang sumber secara tidak berurutan.
  • Simpan secrets di secrets.toml atau env var. Jangan pernah menanam token di kode; pastikan .dlt/secrets.toml tetap di-gitignore.
  • Bekukan skema di produksi pada bagian yang penting. Gunakan schemacontract untuk menangkap perubahan upstream tak terduga alih-alih menyerapnya diam-diam.
  • Periksa lasttrace di CI. Mencatat trace dan jumlah baris memudahkan deteksi kegagalan dan pemuatan nol baris yang senyap.
  • Biarkan dlt menangani EL, dbt menangani T. Tahan godaan menransformasi di dalam resource. Daratkan mentah, transformasi di dbt — ini menjaga lineage bersih dan pemrosesan ulang murah.

Kesimpulan dan Poin Penting

dlt memberi tim data cara native-Python untuk menangani separuh Extract-and-Load dari modern stack tanpa mengadopsi platform konektor yang berat. Inferensi skema otomatis, normalisasi JSON menjadi tabel anak, dan evolusi skema menghilangkan sebagian besar plumbing rapuh yang menumpuk pada skrip ingesti buatan tangan.

Poin penting:

  • dlt (data load tool dari dltHub) adalah pustaka Python untuk EL; ia melengkapi dbt (T) dan Dagster (orkestrasi), serta tidak berhubungan dengan Delta Lake atau PyTorch.
  • Tiga abstraksi inti adalah @dlt.resource, @dlt.source, dan dlt.pipeline dengan sebuah destination.
  • Write disposition (replace, append, merge) ditambah dlt.sources.incremental memberi Anda pemuatan inkremental yang idempoten dengan sebuah primarykey.
  • restapisource deklaratif menangani paginasi, autentikasi, dan penyambungan inkremental dengan konfigurasi alih-alih loop.
  • Secrets berada di .dlt/secrets.toml atau environment variable; schemacontract menegakkan disiplin di produksi.
  • Karena pipeline adalah Python murni, Anda men-deploy-nya dengan cron, GitHub Actions, Airflow, atau Dagster — menyerahkan data mentah ke dbt untuk transformasi.

Mulailah dengan pip install "dlt[duckdb]" dan sebuah API berpaginasi yang sudah Anda kenal; Anda akan punya pipeline yang berfungsi, inkremental, dan sadar skema dalam waktu jauh di bawah satu jam.

Artikel Terkait

Tutorial Dagster: Orkestrasi Data dengan Software-Defined Assets

Dagster: Orkestrasi Data Modern dengan Software-Defined Assets Dagster adalah orkestrator data yang menyusun pipeline be...

Mage: Bikin Data Pipeline Modern yang Rasanya Kayak Ngoding di Notebook

Mage: Bikin Data Pipeline Modern yang Rasanya Kayak Ngoding di Notebook Halo temen-temen, ketemu lagi sama aku Ruby Abdu...

Tutorial Ibis: API DataFrame Portabel untuk Banyak Backend

Ibis: API Dataframe Python yang Portabel di Banyak Backend Ibis adalah library dataframe Python yang memungkinkan Anda m...

Tutorial Pandera: Validasi Data Statistik untuk DataFrame

Pandera: Validasi Data Statistik untuk DataFrame pandas dan Polars Pipeline data sering gagal tanpa suara. Sebuah kolom ...