π Data Engineering β Dari Ingestion ke Data Product
Data engineering adalah disiplin yang mengubah data mentah menjadi aset yang dapat dianalisis . Catatan ini memetakan 7 lapisan pipeline data dari ingestion hingga data product serving, perbandingan batch vs streaming, lakehouse architecture (Iceberg/Delta/Parquet), dan decision framework. Dengan 327 hits di vault, data engineering adalah domain terbanyak kedua yang dirujuk tapi belum memiliki hierarchy master sama sekali.
Daftar Isi
1. Premise β Data Engineering Adalah Tulang Punggung AI
2. Seven-Layer Data Pipeline
3. Layer 0 β Data Source & Ingestion
4. Layer 1 β Storage & Lakehouse
5. Layer 2 β Processing
6. Layer 3 β Transformation (ETL/ELT)
7. Layer 4 β Orchestration & Scheduling
8. Layer 5 β Serving Layer
9. Layer 6 β Observability & Governance
10. Batch vs Streaming
11. Lakehouse Architecture (Iceberg + Delta + Parquet)
1. Premise β Data Engineering Adalah Tulang Punggung AI
Tanpa data engineering yang baik, model AI hanyalah teorema yang tidak bisa diimplementasikan:
[Source Systems] β [Ingestion] β [Storage] β [Processing] β [Serving]
β β β
CDC/Kafka Lakehouse Batch/Stream
Mengapa ini penting:
80% waktu data scientist = data preparation (bukan modeling)
Data pipeline yang buruk β GIGO (Garbage In, Garbage Out)
Data lake tanpa governance β data swamp
Feature store menjembatani engineering dan ML
2. Seven-Layer Data Pipeline
Layer Fungsi Tools Khas Failure Mode L0 Source & Ingestion Kafka, Debezium, Fivetran, Airbyte Schema drift, CDC lag L1 Storage & Lake S3/MinIO, HDFS, Iceberg/Delta/Parquet Corruption, cost explosion L2 Processing Spark, Flink, dbt, Ray Shuffle skew, OOM L3 Transformation dbt, Spark SQL, Airflow Lineage lost, dependency hell L4 Orchestration Airflow, Prefect, Dagster, Temporal DAG failure cascade L5 Serving Trino, Druid, Pinot, ClickHouse Query latency, index miss L6 Observability Datahub, Marquez, Great Expectations Data quality silent failure
3. Layer 0 β Data Source & Ingestion
3.1 Sumber Data
Type Contoh Volume Velocity OLTP DB PostgreSQL, MySQL, SQL Server GB-TB Perubahan per detik SaaS API Salesforce, HubSpot, Stripe MB-GB Per jam Logs Application, server, CDN GB-TB/s Stream IoT Sensor, device telemetry TB-PB High frequency Clickstream Web/mobile events TB-PB Real-time
3.2 CDC (Change Data Capture)
Teknik capture perubahan database tanpa query polling:
CDC Type Mekanisme Latency Overhead Log-based Baca WAL (Write-Ahead Log) ms Minimal Trigger-based DB trigger β audit table ms Signifikan Query-based WHERE updated_at > last_polls-min Tinggi XMIN (Postgres)System column XMIN s Rendah
Tool: Debezium (Kafka Connect), AWS DMS, Fivetran, Airbyte
4. Layer 1 β Storage & Lakehouse
Format Compression Schema Evolution Splittable Use Case CSV/JSON β β β
Ad-hoc, small data Avro β
Moderate β
β
Row-oriented, Kafka Parquet β
Excellent β
β
Analytics (columnar)ORC β
Excellent β
β
Hive/Spark (columnar) Iceberg β
(Parquet) β
(Full DDL) β
β
Lakehouse Delta Lake β
(Parquet) β
(Full DDL) β
β
Lakehouse
Fitur Apache Iceberg Delta Lake Open source β
Apache β LF AI (but open) Time travel β
Snapshot isolation β
Version log ACID β
Row-level β
Row-level Partition evolution β
(hidden partitioning) β (must rewrite) Catalog REST, Hive, Glue, Nessie Unity Catalog, Hive Engine support Spark, Flink, Trino, Presto, Dremio, Snowflake, Athena Spark, Flink, Trino, Presto, Databricks Z-ordering β
Sort order β
Z-order by column Merge β
MERGE INTO β
MERGE INTO
4.3 Lakehouse Architecture
βββββββββββββββββββββββββββββββββββββββββββββββββββββ
β [Serving Layer] Trino, Spark, Snowflake, Athena β
βββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β [Table Format] Iceberg / Delta Lake β
β ββ Snapshot isolation (time travel) β
β ββ ACID transactions β
β ββ Schema evolution β
βββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β [Storage] Object Store (S3, MinIO, GCS, HDFS) β
β ββ Parquet columnar data β
β ββ Manifest / metadata files β
β ββ Compaction (optimize small files) β
βββββββββββββββββββββββββββββββββββββββββββββββββββββ
5. Layer 2 β Processing: Batch & Stream
5.1 Batch Processing
Framework Model Skala Latensi Spark DAG (in-memory) 1-1000 node menit-jam MapReduce (legacy)Disk-based masif jam Hive/Spark SQL SQL declarative 1-1000 node menit-jam dbt SQL + templating Single node menit Pandas/Polars Single-node < 1 node detik-menit
5.2 Stream Processing
Framework Model Semantics State? Apache Flink True streaming Exactly-once β
RocksDB Kafka Streams Library (embedded) Exactly-once β
Local Spark Streaming Micro-batch At-least-once β
State store RisingWave Streaming SQL Exactly-once β
Internal Materalize Streaming SQL Exactly-once β
DuckDB ksqlDB Streaming SQL At-least-once β
Kafka
Benchmark:
Technology Latency Throughput
Flink <100ms 1M+ events/s
Kafka Streams <10ms 500K events/s
Spark Streaming 1-10s 1M+ events/s
RisingWave <1s 100K events/s
6.1 ETL vs ELT
Aspek ETL (Extract-Transform-Load) ELT (Extract-Load-Transform) Transform location Transform engine (Spark) Target data warehouse Schema On-write schema On-read schema Raw data Not preserved Preserved Latency to insight Lebih lambat Lebih cepat Tool modern dbt, Spark, Airflow dbt + Snowflake/BigQuery
Arsitektur dbt:
SQL Model (*.sql) β dbt β Compiled SQL β Run on Warehouse
ββ Lineage graph
ββ Data quality tests
ββ Documentation generation
Model tiers (dbt convention):
Tier Nama Konten Staging stg_* Raw β clean types, rename columns Intermediate int_* Joins, aggregations, business logic Marts dim__, fct__ Kimball: dimension + fact tables Metrics metrics.yml Business metrics layer
7. Layer 4 β Orchestration & Scheduling
7.1 Orchestrator Comparison
Tool DAG Definition Scheduler Backend Retry Airflow Python Time-based Celery/K8s β
Prefect Python Events + time Serverless β
Dagster Python Assets + time K8s β
Temporal Code (Go/Java/Python) Timer + events DB β
Kestra YAML Events + time K8s β
Digdag YAML Time-based DB β
7.2 Airflow DAG Internals
DAG: data_pipeline_v2
βββ extract_from_postgres (PostgresOperator)
β βββ load_raw_to_s3 (S3UploadOperator)
βββ process_with_spark (SparkSubmitOperator)
β βββ [sensor] await_spark_completion
βββ dbt_run_staging (BashOperator)
βββ dbt_run_marts (BashOperator) β depends on dbt_run_staging
βββ dbt_test (BashOperator)
8. Layer 5 β Serving Layer
8.1 OLAP Query Engines
Engine Architecture Query Latency Concurrency Index Trino/Presto Distributed (MPP) 1-30s High β ClickHouse Columnar (MPP) <1s Very High β
Druid Pre-aggregation <1s High β
Pinot Pre-aggregation <100ms Very High β
Snowflake Cloud MPP 1-30s High β BigQuery Serverless 1-30s Unlimited β
8.2 Feature Store
Konsep kunci untuk ML β menjembatani data engineering dan ML engineering:
Feature Store Offline Store Online Store Point-in-time Feast Spark, BigQuery Redis β
Tecton Snowflake, Spark DynamoDB β
Hopsworks HopsFS MySQL Cluster β
Vertex AI Feature Store BigQuery Online store β
Online vs Offline:
Offline: Training dataset β big batch, historical features
Online: Model inference β low-latency, latest feature values
9. Layer 6 β Observability & Governance
9.1 Data Quality
Tool Checks Freshness Schema Great Expectations β
Expectations β β
dbt tests β
Built in β β
Soda β
SQL checks β
β
Deequ (AWS)β
Scala/Spark β β
9.2 Data Catalog
Tool Lineage Discovery Governance Datahub β
β
β
Tags, glossary Amundsen (Lyft)β β
β Apache Atlas β
β
β
Marquez β
β β dbt Docs β
Lineage graph β
β
10. Batch vs Streaming: Decision Framework
10.1 Pilih Berdasarkan
Latency need?
ββ Real-time (< 1s) β Stream processing (Flink/Kafka Streams)
ββ Near real-time (< 1m) β Micro-batch (Spark Streaming)
ββ Minutes β Batch (Spark, dbt)
ββ Hours/days β Scheduled batch
Data volume?
ββ > 100 TB/day β Streaming (Kafka + Flink) atau Incremental batch
ββ 1-100 TB/day β Spark batch
ββ < 1 TB/day β dbt + warehouse
Source pattern?
ββ Continuous (logs, events) β Stream
ββ Scheduled (ERP, nightly) β Batch
ββ CDC (DB changes) β Stream via Debezium
10.2 Lambda vs Kappa Architecture
Aspek Lambda Kappa Batch path β
Seperate β Tidak ada Stream path β
Real-time β
All-stream Complexity 2Γ (batch + stream) 1Γ (stream only) Consistency Reconciliation needed Single source Use case Legacy, batch-heavy Greenfield, stream-native
11. Cross-Reference ke Vault
References
Kleppmann, M. βDesigning Data-Intensive Applications.β OβReilly, 2017.
Kimball, R. & Ross, M. βThe Data Warehouse Toolkit.β 3rd ed., Wiley, 2013.
Narkhede, N., Shapira, G., Palino, T. βKafka: The Definitive Guide.β OβReilly, 2017.
Armbrust, M. et al. βLakehouse: A New Generation of Open Platforms.β CIDR 2021.
Apache Iceberg. βIceberg Table Spec v3.β (2024).
Delta Lake. βDelta Lake Protocol.β (2024).
Carbone, P. et al. βApache Flink: Stream and Batch Processing.β VLDB 2015.
Zaharia, M. et al. βApache Spark: A Unified Engine for Big Data Processing.β CACM 2016.
dbt Labs. βdbt Documentation.β (2024).
Marx, R. βThe Data Engineering Cookbook.β 2020.
Hueske, F. & Kalavri, V. βStream Processing with Apache Flink.β OβReilly, 2019.