πŸ“Š 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. 1. Premise β€” Data Engineering Adalah Tulang Punggung AI
  2. 2. Seven-Layer Data Pipeline
  3. 3. Layer 0 β€” Data Source & Ingestion
  4. 4. Layer 1 β€” Storage & Lakehouse
  5. 5. Layer 2 β€” Processing
  6. 6. Layer 3 β€” Transformation (ETL/ELT)
  7. 7. Layer 4 β€” Orchestration & Scheduling
  8. 8. Layer 5 β€” Serving Layer
  9. 9. Layer 6 β€” Observability & Governance
  10. 10. Batch vs Streaming
  11. 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

LayerFungsiTools KhasFailure Mode
L0Source & IngestionKafka, Debezium, Fivetran, AirbyteSchema drift, CDC lag
L1Storage & LakeS3/MinIO, HDFS, Iceberg/Delta/ParquetCorruption, cost explosion
L2ProcessingSpark, Flink, dbt, RayShuffle skew, OOM
L3Transformationdbt, Spark SQL, AirflowLineage lost, dependency hell
L4OrchestrationAirflow, Prefect, Dagster, TemporalDAG failure cascade
L5ServingTrino, Druid, Pinot, ClickHouseQuery latency, index miss
L6ObservabilityDatahub, Marquez, Great ExpectationsData quality silent failure

3. Layer 0 β€” Data Source & Ingestion

3.1 Sumber Data

TypeContohVolumeVelocity
OLTP DBPostgreSQL, MySQL, SQL ServerGB-TBPerubahan per detik
SaaS APISalesforce, HubSpot, StripeMB-GBPer jam
LogsApplication, server, CDNGB-TB/sStream
IoTSensor, device telemetryTB-PBHigh frequency
ClickstreamWeb/mobile eventsTB-PBReal-time

3.2 CDC (Change Data Capture)

Teknik capture perubahan database tanpa query polling:

CDC TypeMekanismeLatencyOverhead
Log-basedBaca WAL (Write-Ahead Log)msMinimal
Trigger-basedDB trigger β†’ audit tablemsSignifikan
Query-basedWHERE updated_at > last_polls-minTinggi
XMIN (Postgres)System column XMINsRendah

Tool: Debezium (Kafka Connect), AWS DMS, Fivetran, Airbyte


4. Layer 1 β€” Storage & Lakehouse

4.1 File Format Comparison

FormatCompressionSchema EvolutionSplittableUse 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

4.2 Lakehouse Table Format β€” Iceberg vs Delta

FiturApache IcebergDelta 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)
CatalogREST, Hive, Glue, NessieUnity Catalog, Hive
Engine supportSpark, Flink, Trino, Presto, Dremio, Snowflake, AthenaSpark, 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

FrameworkModelSkalaLatensi
SparkDAG (in-memory)1-1000 nodemenit-jam
MapReduce (legacy)Disk-basedmasifjam
Hive/Spark SQLSQL declarative1-1000 nodemenit-jam
dbtSQL + templatingSingle nodemenit
Pandas/PolarsSingle-node< 1 nodedetik-menit

5.2 Stream Processing

FrameworkModelSemanticsState?
Apache FlinkTrue streamingExactly-onceβœ… RocksDB
Kafka StreamsLibrary (embedded)Exactly-onceβœ… Local
Spark StreamingMicro-batchAt-least-onceβœ… State store
RisingWaveStreaming SQLExactly-onceβœ… Internal
MateralizeStreaming SQLExactly-onceβœ… DuckDB
ksqlDBStreaming SQLAt-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. Layer 3 β€” Transformation (ETL/ELT)

6.1 ETL vs ELT

AspekETL (Extract-Transform-Load)ELT (Extract-Load-Transform)
Transform locationTransform engine (Spark)Target data warehouse
SchemaOn-write schemaOn-read schema
Raw dataNot preservedPreserved
Latency to insightLebih lambatLebih cepat
Tool moderndbt, Spark, Airflowdbt + Snowflake/BigQuery

6.2 dbt (Data Build Tool)

Arsitektur dbt:

SQL Model (*.sql) β†’ dbt β†’ Compiled SQL β†’ Run on Warehouse
                           β”œβ”€ Lineage graph
                           β”œβ”€ Data quality tests
                           └─ Documentation generation

Model tiers (dbt convention):

TierNamaKonten
Stagingstg_*Raw β†’ clean types, rename columns
Intermediateint_*Joins, aggregations, business logic
Martsdim__, fct__Kimball: dimension + fact tables
Metricsmetrics.ymlBusiness metrics layer

7. Layer 4 β€” Orchestration & Scheduling

7.1 Orchestrator Comparison

ToolDAG DefinitionSchedulerBackendRetry
AirflowPythonTime-basedCelery/K8sβœ…
PrefectPythonEvents + timeServerlessβœ…
DagsterPythonAssets + timeK8sβœ…
TemporalCode (Go/Java/Python)Timer + eventsDBβœ…
KestraYAMLEvents + timeK8sβœ…
DigdagYAMLTime-basedDBβœ…

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

EngineArchitectureQuery LatencyConcurrencyIndex
Trino/PrestoDistributed (MPP)1-30sHigh❌
ClickHouseColumnar (MPP)<1sVery Highβœ…
DruidPre-aggregation<1sHighβœ…
PinotPre-aggregation<100msVery Highβœ…
SnowflakeCloud MPP1-30sHigh❌
BigQueryServerless1-30sUnlimited❌

8.2 Feature Store

Konsep kunci untuk ML β€” menjembatani data engineering dan ML engineering:

Feature StoreOffline StoreOnline StorePoint-in-time
FeastSpark, BigQueryRedisβœ…
TectonSnowflake, SparkDynamoDBβœ…
HopsworksHopsFSMySQL Clusterβœ…
Vertex AI Feature StoreBigQueryOnline 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

ToolChecksFreshnessSchema
Great Expectationsβœ… ExpectationsβŒβœ…
dbt testsβœ… Built inβŒβœ…
Sodaβœ… SQL checksβœ…βœ…
Deequ (AWS)βœ… Scala/Spark❌❌

9.2 Data Catalog

ToolLineageDiscoveryGovernance
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

AspekLambdaKappa
Batch pathβœ… Seperate❌ Tidak ada
Stream pathβœ… Real-timeβœ… All-stream
Complexity2Γ— (batch + stream)1Γ— (stream only)
ConsistencyReconciliation neededSingle source
Use caseLegacy, batch-heavyGreenfield, stream-native

11. Cross-Reference ke Vault

LayerCatatan Vault
L0hierarchy-abstraction-layers β€” Data layer L8
L1hierarchy-database-storage-systems, hierarchy-memory-storage β€” Tier 5-7
L2hierarchy-concurrency-consensus β€” Distributed processing
L3advanced-chunking-strategies-deepdive β€” Document chunking
L4hierarchy-devops-cicd β€” Pipeline orchestration
L5hierarchy-llm-ai-systems β€” ML serving + RAG
L6hierarchy-cybersecurity-defense-architecture β€” Data governance L3

References

  1. Kleppmann, M. β€œDesigning Data-Intensive Applications.” O’Reilly, 2017.
  2. Kimball, R. & Ross, M. β€œThe Data Warehouse Toolkit.” 3rd ed., Wiley, 2013.
  3. Narkhede, N., Shapira, G., Palino, T. β€œKafka: The Definitive Guide.” O’Reilly, 2017.
  4. Armbrust, M. et al. β€œLakehouse: A New Generation of Open Platforms.” CIDR 2021.
  5. Apache Iceberg. β€œIceberg Table Spec v3.” (2024).
  6. Delta Lake. β€œDelta Lake Protocol.” (2024).
  7. Carbone, P. et al. β€œApache Flink: Stream and Batch Processing.” VLDB 2015.
  8. Zaharia, M. et al. β€œApache Spark: A Unified Engine for Big Data Processing.” CACM 2016.
  9. dbt Labs. β€œdbt Documentation.” (2024).
  10. Marx, R. β€œThe Data Engineering Cookbook.” 2020.
  11. Hueske, F. & Kalavri, V. β€œStream Processing with Apache Flink.” O’Reilly, 2019.