Apache Beam vs Spark 2026: āđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļš Pipeline āđāļšāļšāļĢāļ§āļĄāđāļĨāļ°āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒ

āļ„āļđāđˆāļĄāļ·āļ­āļ„āļĢāļšāļ–āđ‰āļ§āļ™āđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļš Apache Beam 2.76 āđāļĨāļ° Spark 4.2 āļŠāļģāļŦāļĢāļąāļš data engineering āđ€āļĢāļĩāļĒāļ™āļĢāļđāđ‰āļ„āļ§āļēāļĄāđāļ•āļāļ•āđˆāļēāļ‡āļ”āđ‰āļēāļ™āļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄ windowing āļ›āļĢāļ°āļŠāļīāļ—āļ˜āļīāļ āļēāļž āđāļĨāļ°āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļ—āļĩāđˆāļžāļšāļšāđˆāļ­āļĒ

Apache Beam vs Spark 2026 āđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļšāļŠāļģāļŦāļĢāļąāļš data engineering

Apache Beam āđāļĨāļ° Spark āđ€āļ›āđ‡āļ™āļŠāļ­āļ‡ framework āļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨāļ‚āđ‰āļ­āļĄāļđāļĨāļ—āļĩāđˆāđ„āļ”āđ‰āļĢāļąāļšāļ„āļ§āļēāļĄāļ™āļīāļĒāļĄāļŠāļđāļ‡āļŠāļļāļ”āđƒāļ™āļ›āļĩ 2026 Beam 2.76 (āļŠāļīāļ‡āļŦāļēāļ„āļĄ 2026) āđāļĨāļ° Spark 4.2 (āļāļĢāļāļŽāļēāļ„āļĄ 2026) āļ—āļąāđ‰āļ‡āļ„āļđāđˆāļˆāļąāļ”āļāļēāļĢ workload āđāļšāļš batch āđāļĨāļ° streaming āđ„āļ”āđ‰ āđāļ•āđˆāļ›āļĢāļąāļŠāļāļēāļāļēāļĢāļ­āļ­āļāđāļšāļšāļ•āđˆāļēāļ‡āļāļąāļ™āđ‚āļ”āļĒāļžāļ·āđ‰āļ™āļāļēāļ™ Beam āļ—āļģ abstraction āļ‚āļ­āļ‡ execution engine āđƒāļ™āļ‚āļ“āļ°āļ—āļĩāđˆ Spark āļĄāļĩ runtime āđāļšāļšāļĢāļ§āļĄāļ—āļĩāđˆāļœāļŠāļēāļ™āļĢāļ§āļĄāļāļąāļ™āļ­āļĒāđˆāļēāļ‡āđāļ™āđˆāļ™āļŦāļ™āļē

āļāļĢāļ­āļšāļāļēāļĢāļ•āļąāļ”āļŠāļīāļ™āđƒāļˆāļ”āđˆāļ§āļ™

āđ€āļĨāļ·āļ­āļ Beam āđ€āļĄāļ·āđˆāļ­āļ„āļ§āļēāļĄāļŠāļēāļĄāļēāļĢāļ–āđƒāļ™āļāļēāļĢāļžāļāļžāļēāļ‚āđ‰āļēāļĄ runner (Dataflow, Flink, Spark) āļŠāļģāļ„āļąāļāļŦāļĢāļ·āļ­āđ€āļĄāļ·āđˆāļ­āđƒāļŠāđ‰ Google Cloud Dataflow āđ€āļĨāļ·āļ­āļ Spark āđ€āļĄāļ·āđˆāļ­āļˆāļąāļ”āļāļēāļĢ cluster āđ€āļ­āļ‡ āđƒāļŠāđ‰ Databricks āļŦāļĢāļ·āļ­āļ•āđ‰āļ­āļ‡āļāļēāļĢāļāļēāļĢāļœāļŠāļēāļ™āļĢāļ§āļĄ ML āļāļąāļš MLlib

āđ‚āļĄāđ€āļ”āļĨāļ„āļ§āļēāļĄāļžāļāļžāļēāļ‚āļ­āļ‡ Beam vs Engine āđāļšāļšāļĢāļ§āļĄāļ‚āļ­āļ‡ Spark

Apache Beam āđāļĒāļāđ‚āļĄāđ€āļ”āļĨāļāļēāļĢāđ€āļ‚āļĩāļĒāļ™āđ‚āļ›āļĢāđāļāļĢāļĄāļ­āļ­āļāļˆāļēāļāļāļēāļĢāļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨ āļ„āļģāļˆāļģāļāļąāļ”āļ„āļ§āļēāļĄ pipeline āđ€āļ”āļĩāļĒāļ§āļŠāļēāļĄāļēāļĢāļ–āļ—āļģāļ‡āļēāļ™āļšāļ™ Google Cloud Dataflow Apache Flink Apache Spark āļŦāļĢāļ·āļ­ runner āļ­āļ·āđˆāļ™āđ„āļ”āđ‰āđ‚āļ”āļĒāđ„āļĄāđˆāļ•āđ‰āļ­āļ‡āđ€āļ›āļĨāļĩāđˆāļĒāļ™āđ‚āļ„āđ‰āļ” āļāļēāļĢ abstraction āļ™āļĩāđ‰āļĄāļēāļˆāļēāļ Beam SDK āļ—āļĩāđˆāļŠāļĢāđ‰āļēāļ‡āļāļēāļĢāđāļŠāļ”āļ‡āļœāļĨ pipeline āđāļšāļšāļžāļāļžāļēāļ—āļĩāđˆ runner āļ—āļĩāđˆāđ€āļ‚āđ‰āļēāļāļąāļ™āđ„āļ”āđ‰āļŠāļēāļĄāļēāļĢāļ–āļ•āļĩāļ„āļ§āļēāļĄāđ„āļ”āđ‰

Spark 4.2 āđƒāļŠāđ‰āđāļ™āļ§āļ—āļēāļ‡āļ•āļĢāļ‡āļāļąāļ™āļ‚āđ‰āļēāļĄ DataFrame API, Structured Streaming āđāļĨāļ° MLlib āđƒāļŠāđ‰ Catalyst optimizer āđāļĨāļ° Tungsten execution engine āļĢāđˆāļ§āļĄāļāļąāļ™ āļāļēāļĢāļœāļŠāļēāļ™āļĢāļ§āļĄāļ­āļĒāđˆāļēāļ‡āđāļ™āđˆāļ™āļŦāļ™āļēāļ™āļĩāđ‰āļ—āļģāđƒāļŦāđ‰āļŠāļēāļĄāļēāļĢāļ–āļ›āļĢāļąāļšāđāļ•āđˆāļ‡āļ­āļĒāđˆāļēāļ‡ Adaptive Query Execution āļ—āļĩāđˆāļ›āļĢāļąāļšāđāļœāļ™āļ‚āļ“āļ° runtime āļ•āļēāļĄāļŠāļ–āļīāļ•āļīāļ‚āđ‰āļ­āļĄāļđāļĨāļˆāļĢāļīāļ‡

python
# beam_pipeline.py
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

# āđ‚āļ„āđ‰āļ”āđ€āļ”āļĩāļĒāļ§āļāļąāļ™āļ—āļģāļ‡āļēāļ™āļšāļ™ Dataflow, Flink āļŦāļĢāļ·āļ­ Spark runner
options = PipelineOptions([
    '--runner=DataflowRunner',  # āļŠāļĨāļąāļšāđ€āļ›āđ‡āļ™ FlinkRunner āļŦāļĢāļ·āļ­ SparkRunner
    '--project=my-project',
    '--region=us-central1',
    '--temp_location=gs://my-bucket/temp'
])

with beam.Pipeline(options=options) as pipeline:
    (pipeline
     | 'ReadEvents' >> beam.io.ReadFromPubSub(topic='projects/p/topics/events')
     | 'ParseJSON' >> beam.Map(lambda x: json.loads(x))
     | 'FilterValid' >> beam.Filter(lambda e: e.get('status') == 'valid')
     | 'WindowByMinute' >> beam.WindowInto(beam.window.FixedWindows(60))
     | 'CountPerWindow' >> beam.combiners.Count.Globally()
     | 'WriteToBQ' >> beam.io.WriteToBigQuery('project:dataset.table'))

āđ€āļ§āļ­āļĢāđŒāļŠāļąāļ™ Spark āļ—āļĩāđˆāđ€āļ—āļĩāļĒāļšāđ€āļ—āđˆāļēāļœāļđāļāļ•āļīāļ”āļāļąāļš Spark runtime āđ‚āļ”āļĒāļ•āļĢāļ‡:

python
# spark_streaming.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, window

spark = SparkSession.builder \
    .appName("EventProcessing") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

# Spark 4.2: āđ‚āļŦāļĄāļ” ANSI āđ€āļ›āļīāļ”āđƒāļŠāđ‰āļ‡āļēāļ™āđ‚āļ”āļĒāļ„āđˆāļēāđ€āļĢāļīāđˆāļĄāļ•āđ‰āļ™ āļāļēāļĢāļ•āļĢāļ§āļˆāļŠāļ­āļšāļ›āļĢāļ°āđ€āļ āļ—āđ€āļ‚āđ‰āļĄāļ‡āļ§āļ”āļ‚āļķāđ‰āļ™
events = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "broker:9092") \
    .option("subscribe", "events") \
    .load()

processed = events \
    .select(from_json(col("value").cast("string"), schema).alias("data")) \
    .filter(col("data.status") == "valid") \
    .groupBy(window(col("data.timestamp"), "1 minute")) \
    .count()

processed.writeStream \
    .format("bigquery") \
    .option("table", "project.dataset.table") \
    .outputMode("append") \
    .start()

āļ„āļ§āļēāļĄāļžāļāļžāļēāļ‚āļ­āļ‡ Beam āļ—āļģāđƒāļŦāđ‰āļŠāļēāļĄāļēāļĢāļ–āļŠāļĢāđ‰āļēāļ‡āļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄ cloud-agnostic āđ„āļ”āđ‰āđāļ•āđˆāđ€āļžāļīāđˆāļĄ translation layer āļāļēāļĢāļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨāđ‚āļ”āļĒāļ•āļĢāļ‡āļ‚āļ­āļ‡ Spark āđ‚āļ”āļĒāļ—āļąāđˆāļ§āđ„āļ›āđāļŠāļ”āļ‡ latency āļ•āđˆāļģāļāļ§āđˆāļēāļŠāļģāļŦāļĢāļąāļšāļāļēāļĢāļ”āļģāđ€āļ™āļīāļ™āļāļēāļĢāļ—āļĩāđˆāđ€āļ—āļĩāļĒāļšāđ€āļ—āđˆāļēāļšāļ™ hardware āđ€āļ”āļĩāļĒāļ§āļāļąāļ™

āđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļš Windowing āđāļĨāļ°āļāļēāļĢāļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨ Event-Time

āļ—āļąāđ‰āļ‡āļŠāļ­āļ‡ framework āļˆāļąāļ”āļāļēāļĢ semantics āļ‚āļ­āļ‡ event-time āđāļ•āđˆ API āļŠāļ°āļ—āđ‰āļ­āļ™āļĄāļĢāļ”āļāļ—āļĩāđˆāļ•āđˆāļēāļ‡āļāļąāļ™ āđ‚āļĄāđ€āļ”āļĨ windowing āļ‚āļ­āļ‡ Beam āļĄāļēāļˆāļēāļ Dataflow Model paper (2015) āļ‹āļķāđˆāļ‡āļ–āļ·āļ­āļ§āđˆāļē window āđ€āļ›āđ‡āļ™āļ­āļ‡āļ„āđŒāļ›āļĢāļ°āļāļ­āļš pipeline āļŠāļąāđ‰āļ™āļ™āļģ Spark āļ›āļĢāļąāļš windowing āļŠāļģāļŦāļĢāļąāļš Structured Streaming āđ‚āļ”āļĒāļœāļŠāļēāļ™āļĢāļ§āļĄāļāļąāļš DataFrame API

āļ„āļļāļ“āļŠāļĄāļšāļąāļ•āļīBeam 2.76Spark 4.2
Fixed windowsFixedWindows(duration)window(col, duration)
Sliding windowsSlidingWindows(size, period)window(col, size, period)
Session windowsSessions(gap)āđ„āļĄāđˆ native (āđƒāļŠāđ‰ flatMapGroupsWithState)
Custom windowsSubclass WindowFnāļˆāļģāļāļąāļ”
āļˆāļąāļ”āļāļēāļĢ late dataTrigger āđƒāļ™āļ•āļąāļ§Watermark delays
Allowed latenessāļāļģāļŦāļ™āļ”āļ„āđˆāļēāļ•āđˆāļ­ windowGlobal watermark

Session windows āđāļŠāļ”āļ‡āļ„āļ§āļēāļĄāđāļ•āļāļ•āđˆāļēāļ‡āļŠāļąāļ”āđ€āļˆāļ™āļ—āļĩāđˆāļŠāļļāļ” Beam āļ–āļ·āļ­āļ§āđˆāļē session āđ€āļ›āđ‡āļ™āļāļĨāļĒāļļāļ—āļ˜āđŒ windowing āđāļšāļš native:

python
# beam_sessions.py
from apache_beam import window

# Session gap 30 āļ™āļēāļ—āļĩ āļĢāļąāļšāļ‚āđ‰āļ­āļĄāļđāļĨāļĨāđˆāļēāļŠāđ‰āļēāļ–āļķāļ‡ 1 āļŠāļąāđˆāļ§āđ‚āļĄāļ‡
windowed = (
    events
    | 'SessionWindow' >> beam.WindowInto(
        window.Sessions(30 * 60),  # Gap 30 āļ™āļēāļ—āļĩāļ›āļīāļ” session
        trigger=beam.trigger.AfterWatermark(
            early=beam.trigger.AfterProcessingTime(60),
            late=beam.trigger.AfterCount(1)
        ),
        allowed_lateness=3600,  # āļĢāļąāļšāļ‚āđ‰āļ­āļĄāļđāļĨāļĨāđˆāļēāļŠāđ‰āļēāļ–āļķāļ‡ 1 āļŠāļąāđˆāļ§āđ‚āļĄāļ‡
        accumulation_mode=beam.trigger.AccumulationMode.ACCUMULATING
    )
)

Spark āļ•āđ‰āļ­āļ‡āļāļēāļĢāļāļēāļĢāļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨāđāļšāļš stateful āļŠāļģāļŦāļĢāļąāļš sessions:

python
# spark_sessions.py
from pyspark.sql.streaming import GroupState, GroupStateTimeout

def update_session(key, events, state: GroupState):
    # āļˆāļąāļ”āļāļēāļĢ session āļ”āđ‰āļ§āļĒāļ•āļ™āđ€āļ­āļ‡āļāļąāļš state API
    session_data = state.getOption() or {"count": 0, "start": None, "end": None}
    
    for event in events:
        ts = event.timestamp
        if session_data["end"] and (ts - session_data["end"]).seconds > 1800:
            # Gap āđ€āļāļīāļ™ 30 āļ™āļēāļ—āļĩ āļŠāđˆāļ‡āļ­āļ­āļ session āļāđˆāļ­āļ™āļŦāļ™āđ‰āļē
            yield session_data
            session_data = {"count": 0, "start": ts, "end": ts}
        
        session_data["count"] += 1
        session_data["end"] = ts
        if not session_data["start"]:
            session_data["start"] = ts
    
    state.update(session_data)
    state.setTimeoutDuration(30 * 60 * 1000)  # Timeout 30 āļ™āļēāļ—āļĩ

# Spark 4.2 Arbitrary State API v2
result = events \
    .groupByKey(lambda e: e.user_id) \
    .flatMapGroupsWithState(
        update_session,
        outputMode="append",
        stateType=session_schema,
        timeoutConf=GroupStateTimeout.ProcessingTimeTimeout
    )

āļŠāļģāļŦāļĢāļąāļšāļāļēāļĢāđ€āļ•āļĢāļĩāļĒāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāđ€āļāļĩāđˆāļĒāļ§āļāļąāļšāļŦāļąāļ§āļ‚āđ‰āļ­āđ€āļŦāļĨāđˆāļēāļ™āļĩāđ‰ āļ”āļđāđ‚āļĄāļ”āļđāļĨ āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒ Apache Beam āđāļĨāļ° Dataflow

āļ›āļĢāļ°āļŠāļīāļ—āļ˜āļīāļ āļēāļž: āļ‚āđ‰āļ­āļĄāļđāļĨ Benchmark 2026

āļāļēāļĢāđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļšāđ‚āļ”āļĒāļ•āļĢāļ‡āļ•āđ‰āļ­āļ‡āļāļēāļĢāļāļēāļĢāļ•āļąāđ‰āļ‡āļ„āđˆāļēāļ­āļĒāđˆāļēāļ‡āļĢāļ°āļĄāļąāļ”āļĢāļ°āļ§āļąāļ‡āđ€āļžāļĢāļēāļ° Beam āļ—āļģāļ‡āļēāļ™āļšāļ™ Spark āđ€āļ›āđ‡āļ™āļ•āļąāļ§āđ€āļĨāļ·āļ­āļ runner āļŦāļ™āļķāđˆāļ‡ āļāļēāļĢāđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļšāļ—āļĩāđˆāđ€āļāļĩāđˆāļĒāļ§āļ‚āđ‰āļ­āļ‡āļ„āļ·āļ­ Beam-on-Dataflow vs native Spark

Benchmark āļĨāđˆāļēāļŠāļļāļ”āļˆāļēāļ Databricks āđāļĨāļ° Google Cloud āđāļŠāļ”āļ‡:

WorkloadSpark 4.2 (Databricks)Beam 2.76 (Dataflow)āļŦāļĄāļēāļĒāđ€āļŦāļ•āļļ
Batch ETL (1TB Parquet)4.2 āļ™āļēāļ—āļĩ5.1 āļ™āļēāļ—āļĩāļ‚āđ‰āļ­āđ„āļ”āđ‰āđ€āļ›āļĢāļĩāļĒāļš Photon engine āļ‚āļ­āļ‡ Spark
Streaming (100K events/āļ§āļīāļ™āļēāļ—āļĩ)45ms p99 latency120ms p99 latencyOverhead autoscaling āļ‚āļ­āļ‡ Dataflow
Exactly-once sink writesNativeNativeāļ—āļąāđ‰āļ‡āļ„āļđāđˆāļĢāļ­āļ‡āļĢāļąāļšāļ•āļąāđ‰āļ‡āđāļ•āđˆ 2024
āļ„āđˆāļēāđƒāļŠāđ‰āļˆāđˆāļēāļĒ (workload āļ•āđˆāļ­āđ€āļ™āļ·āđˆāļ­āļ‡)$0.12/GB āļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨ$0.08/GB āļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨPricing Dataflow Flex

āļāļēāļĢāļ›āļĢāļąāļšāļ›āļĢāļļāļ‡āļ›āļĢāļ°āļŠāļīāļ—āļ˜āļīāļ āļēāļžāļ‚āļ­āļ‡ Spark 4.2 āļĄāļēāļˆāļēāļāļŦāļĨāļēāļĒāļ„āļļāļ“āļŠāļĄāļšāļąāļ•āļī:

  • āđ‚āļŦāļĄāļ” ANSI āđ€āļĢāļīāđˆāļĄāļ•āđ‰āļ™: semantics SQL āđ€āļ‚āđ‰āļĄāļ‡āļ§āļ”āļ‚āļķāđ‰āļ™āļˆāļąāļš error āđ€āļĢāđ‡āļ§āļ‚āļķāđ‰āļ™
  • āļ›āļĢāļ°āđ€āļ āļ—āļ‚āđ‰āļ­āļĄāļđāļĨ VARIANT: āļˆāļąāļ”āļāļēāļĢāļ‚āđ‰āļ­āļĄāļđāļĨāļāļķāđˆāļ‡āđ‚āļ„āļĢāļ‡āļŠāļĢāđ‰āļēāļ‡āđāļšāļš native
  • Adaptive Query Execution: āļ›āļĢāļąāļšāđāļ•āđˆāļ‡āđāļœāļ™āļ‚āļ“āļ° runtime
  • āļĢāļ­āļ‡āļĢāļąāļš Java 21: virtual threads āļĨāļ” overhead

Beam 2.76 āļ™āļģāđ€āļŠāļ™āļ­:

  • āļĢāļ­āļ‡āļĢāļąāļš Flink 2.0 runner: stateful streaming āļĢāļ°āļ”āļąāļš production
  • CDC offset persistence: DebeziumIO FileSystemOffsetRetainer
  • āļāļēāļĢāļœāļŠāļēāļ™āļĢāļ§āļĄ ADK: āļĢāļ­āļ‡āļĢāļąāļš Google Agent Development Kit āđƒāļ™ Python SDK

āļžāļĢāđ‰āļ­āļĄāļ—āļĩāđˆāļˆāļ°āļžāļīāļŠāļīāļ•āļāļēāļĢāļŠāļąāļĄāļ āļēāļĐāļ“āđŒ Data Engineering āđāļĨāđ‰āļ§āļŦāļĢāļ·āļ­āļĒāļąāļ‡āļ„āļĢāļąāļš?

āļāļķāļāļāļ™āļ”āđ‰āļ§āļĒāļ•āļąāļ§āļˆāļģāļĨāļ­āļ‡āđāļšāļšāđ‚āļ•āđ‰āļ•āļ­āļš, flashcards āđāļĨāļ°āđāļšāļšāļ—āļ”āļŠāļ­āļšāđ€āļ—āļ„āļ™āļīāļ„āļ„āļĢāļąāļš

āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļ—āļĩāđˆāļžāļšāļšāđˆāļ­āļĒ: Beam vs Spark

āļāļēāļĢāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāđ€āļ—āļ„āļ™āļīāļ„āļŠāļģāļŦāļĢāļąāļšāļ•āļģāđāļŦāļ™āđˆāļ‡ data engineering āļĄāļąāļāđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļšāļŠāļ­āļ‡ framework āļ™āļĩāđ‰ āļ™āļĩāđˆāļ„āļ·āļ­āļ„āļģāļ–āļēāļĄāļ—āļĩāđˆāļ›āļĢāļēāļāļāđƒāļ™āļāļēāļĢāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļ›āļĩ 2026 āļžāļĢāđ‰āļ­āļĄāļ„āļ§āļēāļĄāļĨāļķāļāļ‚āļ­āļ‡āļ„āļģāļ•āļ­āļšāļ—āļĩāđˆāļ„āļēāļ”āļŦāļ§āļąāļ‡

Q1: āđ€āļĄāļ·āđˆāļ­āđ„āļŦāļĢāđˆāļ„āļ§āļĢāđ€āļĨāļ·āļ­āļ Beam āđāļ—āļ™ native Spark?

āļ›āļĢāļ°āđ€āļ”āđ‡āļ™āļ„āļģāļ•āļ­āļšāļ—āļĩāđˆāļ„āļēāļ”āļŦāļ§āļąāļ‡:

  • āļāļēāļĢ deploy āđāļšāļš multi-cloud āļŦāļĢāļ·āļ­ hybrid āļ—āļĩāđˆāđ‚āļ„āđ‰āļ” pipeline āļ•āđ‰āļ­āļ‡āļ—āļģāļ‡āļēāļ™āđ‚āļ”āļĒāđ„āļĄāđˆāđ€āļ›āļĨāļĩāđˆāļĒāļ™āđāļ›āļĨāļ‡
  • āļŠāļ āļēāļžāđāļ§āļ”āļĨāđ‰āļ­āļĄ Google Cloud āļ—āļĩāđˆ Dataflow āđƒāļŦāđ‰āđ‚āļ„āļĢāļ‡āļŠāļĢāđ‰āļēāļ‡āļžāļ·āđ‰āļ™āļāļēāļ™āļ—āļĩāđˆāļˆāļąāļ”āļāļēāļĢāđƒāļŦāđ‰
  • Semantics event-time āļ‹āļąāļšāļ‹āđ‰āļ­āļ™āļ—āļĩāđˆāļ•āđ‰āļ­āļ‡āļāļēāļĢ session windows āļŦāļĢāļ·āļ­ custom triggers
  • āļ—āļĩāļĄāļ—āļĩāđˆāļĄāļĩāļ„āļ§āļēāļĄāđ€āļŠāļĩāđˆāļĒāļ§āļŠāļēāļ Beam āļˆāļēāļāļžāļ·āđ‰āļ™āļāļēāļ™ Dataflow

āļ„āļģāļ•āļ­āļšāļ˜āļ‡āđāļ”āļ‡: "Beam āļ”āļĩāļāļ§āđˆāļēāđ€āļŠāļĄāļ­āđ€āļžāļĢāļēāļ°āļžāļāļžāļēāđ„āļ”āđ‰" āļ„āļ§āļēāļĄāļžāļāļžāļēāļĄāļĩāļ„āđˆāļēāđƒāļŠāđ‰āļˆāđˆāļēāļĒ overhead

Q2: āļāļēāļĢ abstraction runner āļ‚āļ­āļ‡ Beam āļŠāđˆāļ‡āļœāļĨāļ•āđˆāļ­āļāļēāļĢ debug āļ­āļĒāđˆāļēāļ‡āđ„āļĢ?

āļ›āļĢāļ°āđ€āļ”āđ‡āļ™āļ„āļģāļ•āļ­āļšāļ—āļĩāđˆāļ„āļēāļ”āļŦāļ§āļąāļ‡:

  • Stack trace āļ­āđ‰āļēāļ‡āļ­āļīāļ‡āļ—āļąāđ‰āļ‡ Beam SDK āđāļĨāļ°āļāļēāļĢ implement āļ‚āļ­āļ‡ runner
  • āļāļēāļĢāļ›āļĢāļąāļšāđāļ•āđˆāļ‡āđ€āļ‰āļžāļēāļ° runner āļ­āļēāļˆāļ—āļģāļ‡āļēāļ™āļ•āđˆāļēāļ‡āļāļąāļ™ (checkpoint Flink vs checkpoint Spark)
  • Metrics API āđƒāļŦāđ‰ monitoring āđāļšāļšāļĢāļ§āļĄ āđāļ•āđˆ dashboard runner āđāļŠāļ”āļ‡āļĢāļēāļĒāļĨāļ°āđ€āļ­āļĩāļĒāļ”āļ•āđˆāļēāļ‡āļāļąāļ™
  • āļ—āļ”āļŠāļ­āļšāļāļąāļš DirectRunner āļāđˆāļ­āļ™ deploy āđ„āļ›āļĒāļąāļ‡ production runner

Q3: āļ­āļ˜āļīāļšāļēāļĒ semantics exactly-once āđƒāļ™āļ—āļąāđ‰āļ‡āļŠāļ­āļ‡ framework

āļ„āļģāļ•āļ­āļšāļ—āļĩāđˆāļ„āļēāļ”āļŦāļ§āļąāļ‡:

python
# Beam: exactly-once āļœāđˆāļēāļ™āļāļēāļĢāļĢāļąāļšāļ›āļĢāļ°āļāļąāļ™āļ‚āļ­āļ‡ runner
# Dataflow āđƒāļŦāđ‰ exactly-once āļŠāļģāļŦāļĢāļąāļšāļ—āļąāđ‰āļ‡ sources āđāļĨāļ° sinks
# SDK āļˆāļąāļ”āļāļēāļĢāļāļēāļĢāļĨāļšāļ‹āđ‰āļģāđāļĨāļ°āļāļēāļĢāļ›āļĢāļ°āļŠāļēāļ™āļ‡āļēāļ™ checkpoint

with beam.Pipeline() as p:
    (p
     | beam.io.ReadFromPubSub(subscription='...')  # āļ­āđˆāļēāļ™ exactly-once
     | beam.Map(process)
     | beam.io.WriteToBigQuery(...)  # āđ€āļ‚āļĩāļĒāļ™ exactly-once āļžāļĢāđ‰āļ­āļĄ retry
    )

# Spark: exactly-once āļœāđˆāļēāļ™ checkpointing āđāļĨāļ° idempotent sinks
spark.readStream \
    .format("kafka") \
    .load() \
    .writeStream \
    .option("checkpointLocation", "/checkpoint")  # āļāļđāđ‰āļ„āļ·āļ™ state
    .foreachBatch(idempotent_write)  # āļĨāļšāļ‹āđ‰āļģāļĢāļ°āļ”āļąāļšāđāļ­āļ›āļžāļĨāļīāđ€āļ„āļŠāļąāļ™
    .start()

āļŠāļģāļŦāļĢāļąāļšāļŦāļąāļ§āļ‚āđ‰āļ­āļŠāļąāļĄāļ āļēāļĐāļ“āđŒ streaming āđ€āļžāļīāđˆāļĄāđ€āļ•āļīāļĄ āļ”āļđ āđ‚āļĄāļ”āļđāļĨ PySpark

Q4: āļˆāļ° migrate Spark batch job āđ„āļ›āļĒāļąāļ‡ Beam āļ­āļĒāđˆāļēāļ‡āđ„āļĢ?

āļ„āļģāļ–āļēāļĄāļ™āļĩāđ‰āļ—āļ”āļŠāļ­āļšāļ„āļ§āļēāļĄāđ€āļ‚āđ‰āļēāđƒāļˆāļ—āļąāđ‰āļ‡āļŠāļ­āļ‡ API āļ›āļĢāļ°āđ€āļ”āđ‡āļ™āļŠāļģāļ„āļąāļ:

  1. Map āļāļēāļĢāļ”āļģāđ€āļ™āļīāļ™āļāļēāļĢ DataFrame āđ„āļ›āļĒāļąāļ‡ PCollections āđāļĨāļ° transforms
  2. āđāļ—āļ™āļ—āļĩāđˆ spark.read āļ”āđ‰āļ§āļĒ Beam I/O connector āļ—āļĩāđˆāđ€āļŦāļĄāļēāļ°āļŠāļĄ
  3. āđāļ›āļĨāļ‡ UDF āđ€āļ›āđ‡āļ™āļŸāļąāļ‡āļāđŒāļŠāļąāļ™ beam.Map āļŦāļĢāļ·āļ­ beam.ParDo
  4. āļˆāļąāļ”āļāļēāļĢ partitioning āļ•āđˆāļēāļ‡āļāļąāļ™ (Reshuffle āļ‚āļ­āļ‡ Beam vs repartition āļ‚āļ­āļ‡ Spark)
  5. āļ—āļ”āļŠāļ­āļšāļāļąāļš DirectRunner āļāđˆāļ­āļ™ deploy āđ„āļ›āļĒāļąāļ‡ production runner

āđ€āļĨāļ·āļ­āļāđ€āļ„āļĢāļ·āđˆāļ­āļ‡āļĄāļ·āļ­āļ—āļĩāđˆāđ€āļŦāļĄāļēāļ°āļŠāļĄ: āđ€āļĄāļ—āļĢāļīāļāļ‹āđŒāļāļēāļĢāļ•āļąāļ”āļŠāļīāļ™āđƒāļˆ

āļāļēāļĢāđ€āļĨāļ·āļ­āļāļ‚āļķāđ‰āļ™āļ­āļĒāļđāđˆāļāļąāļšāļšāļĢāļīāļšāļ—āļ‚āļ­āļ‡āļ­āļ‡āļ„āđŒāļāļĢāļĄāļēāļāļāļ§āđˆāļēāļ„āļ§āļēāļĄāļŠāļēāļĄāļēāļĢāļ–āļ—āļēāļ‡āđ€āļ—āļ„āļ™āļīāļ„ āļ—āļąāđ‰āļ‡āļŠāļ­āļ‡ framework āļˆāļąāļ”āļāļēāļĢ workload data engineering āļŠāđˆāļ§āļ™āđƒāļŦāļāđˆāđ„āļ”āđ‰āļ­āļĒāđˆāļēāļ‡āļĄāļĩāļ›āļĢāļ°āļŠāļīāļ—āļ˜āļīāļ āļēāļž

āļ›āļąāļˆāļˆāļąāļĒāļŠāļ™āļąāļšāļŠāļ™āļļāļ™ BeamāļŠāļ™āļąāļšāļŠāļ™āļļāļ™ Spark
Cloud providerGoogle CloudAWS EMR, Databricks, on-prem
āļ„āļ§āļēāļĄāđ€āļŠāļĩāđˆāļĒāļ§āļŠāļēāļāļ—āļĩāļĄāļ›āļĢāļ°āļŠāļšāļāļēāļĢāļ“āđŒ Dataflowāļ—āļąāļāļĐāļ° Spark/PySpark
āļāļēāļĢāļœāļŠāļēāļ™āļĢāļ§āļĄ MLāļˆāļģāļāļąāļ” (āđ€āļ„āļĢāļ·āđˆāļ­āļ‡āļĄāļ·āļ­āđāļĒāļ)MLlib, Spark ML
āļāļēāļĢāļ§āļīāđ€āļ„āļĢāļēāļ°āļŦāđŒāđāļšāļšāđ‚āļ•āđ‰āļ•āļ­āļšāđ„āļĄāđˆāđ„āļ”āđ‰āļ­āļ­āļāđāļšāļšāļĄāļēāļŠāļģāļŦāļĢāļąāļšāļŠāļīāđˆāļ‡āļ™āļĩāđ‰Spark SQL, notebooks
Session windowsāļĢāļ­āļ‡āļĢāļąāļš nativeāļˆāļąāļ”āļāļēāļĢ state āļ”āđ‰āļ§āļĒāļ•āļ™āđ€āļ­āļ‡
āđ‚āļĄāđ€āļ”āļĨāļ„āđˆāļēāđƒāļŠāđ‰āļˆāđˆāļēāļĒPay-per-use (Dataflow)āļˆāļąāļ”āļŠāļĢāļĢ cluster
Vendor lock-ināļ•āđˆāļģāļāļ§āđˆāļē (āļŦāļĨāļēāļĒ runners)āļŠāļđāļ‡āļāļ§āđˆāļē (āđ‚āļ„āđ‰āļ”āđ€āļ‰āļžāļēāļ° Spark)

āļŠāļģāļŦāļĢāļąāļš āļāļēāļĢāļ•āļąāļ”āļŠāļīāļ™āđƒāļˆāļĢāļđāļ›āđāļšāļš ETL/ELT āļžāļīāļˆāļēāļĢāļ“āļēāļ›āļĢāļīāļĄāļēāļ“āļ‚āđ‰āļ­āļĄāļđāļĨāđāļĨāļ°āļ‚āđ‰āļ­āļāļģāļŦāļ™āļ” latency āļāđˆāļ­āļ™āđ€āļĨāļ·āļ­āļ framework

āļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄāđƒāļ™āđ‚āļĨāļāļˆāļĢāļīāļ‡: āđāļ™āļ§āļ—āļēāļ‡ Hybrid

āļŦāļĨāļēāļĒāļ­āļ‡āļ„āđŒāļāļĢāđƒāļŠāđ‰āļ—āļąāđ‰āļ‡āļŠāļ­āļ‡ framework āļĢāļđāļ›āđāļšāļšāļ—āļąāđˆāļ§āđ„āļ›:

text
[Streaming Ingestion]     [Batch Processing]     [ML Training]
        |                        |                    |
   Beam/Dataflow           Spark āļšāļ™ Databricks    Spark MLlib
        |                        |                    |
        v                        v                    v
    BigQuery  <----- dbt -----> Delta Lake -----> Model Registry

Beam āļˆāļąāļ”āļāļēāļĢ streaming ingestion āļ—āļĩāđˆ autoscaling āļ‚āļ­āļ‡ Dataflow āļ•āļĢāļ‡āļāļąāļšāļĢāļđāļ›āđāļšāļš traffic Spark āļ›āļĢāļ°āļĄāļ§āļĨāļœāļĨ batch workload āļ—āļĩāđˆāđ€āļĻāļĢāļĐāļāļĻāļēāļŠāļ•āļĢāđŒāļ‚āļ­āļ‡ cluster āļ”āļĩāļāļ§āđˆāļēāļŠāļģāļŦāļĢāļąāļšāļāļēāļĢāļ„āļģāļ™āļ§āļ“āļ•āđˆāļ­āđ€āļ™āļ·āđˆāļ­āļ‡ āļ—āļąāđ‰āļ‡āļ„āļđāđˆāđ„āļŦāļĨāđ„āļ›āļĒāļąāļ‡ data warehouse āđāļšāļšāļĢāļ§āļĄ

āļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄāļ™āļĩāđ‰āļāļĨāđˆāļēāļ§āļ–āļķāļ‡āđƒāļ™ āļšāļ—āđ€āļĢāļĩāļĒāļ™ Apache Spark āļžāļĢāđ‰āļ­āļĄāļĢāļēāļĒāļĨāļ°āđ€āļ­āļĩāļĒāļ”āļāļēāļĢāļ™āļģāđ„āļ›āđƒāļŠāđ‰

āđ€āļĢāļīāđˆāļĄāļāļķāļāļ‹āđ‰āļ­āļĄāđ€āļĨāļĒ!

āļ—āļ”āļŠāļ­āļšāļ„āļ§āļēāļĄāļĢāļđāđ‰āļ‚āļ­āļ‡āļ„āļļāļ“āļ”āđ‰āļ§āļĒāļ•āļąāļ§āļˆāļģāļĨāļ­āļ‡āļŠāļąāļĄāļ āļēāļĐāļ“āđŒāđāļĨāļ°āđāļšāļšāļ—āļ”āļŠāļ­āļšāđ€āļ—āļ„āļ™āļīāļ„āļ„āļĢāļąāļš

āļŠāļĢāļļāļ›āļāļēāļĢāđ€āļĨāļ·āļ­āļ Beam vs Spark

  • Beam 2.76 āđƒāļŦāđ‰āļ„āļ§āļēāļĄāļŠāļēāļĄāļēāļĢāļ–āđƒāļ™āļāļēāļĢāļžāļāļžāļē runner āļšāļ™ Dataflow, Flink 2.0 āđāļĨāļ° Spark āđāļĨāļāđ€āļ›āļĨāļĩāđˆāļĒāļ™āļ›āļĢāļ°āļŠāļīāļ—āļ˜āļīāļ āļēāļžāđ€āļžāļ·āđˆāļ­āļ„āļ§āļēāļĄāļĒāļ·āļ”āļŦāļĒāļļāđˆāļ™āđƒāļ™āļāļēāļĢ deploy
  • Spark 4.2 āđƒāļŦāđ‰āļāļēāļĢāļœāļŠāļēāļ™āļĢāļ§āļĄāļ—āļĩāđˆāđāļ™āđˆāļ™āļ‚āļķāđ‰āļ™āļĢāļ°āļŦāļ§āđˆāļēāļ‡ SQL, streaming āđāļĨāļ° ML workload āļ”āđ‰āļ§āļĒāļ„āļļāļ“āļŠāļĄāļšāļąāļ•āļīāđ€āļŠāđˆāļ™āļ›āļĢāļ°āđ€āļ āļ— VARIANT āđāļĨāļ° Adaptive Query Execution
  • Session windows āđāļĨāļ° trigger āļ‹āļąāļšāļ‹āđ‰āļ­āļ™āđ€āļŦāļĄāļēāļ°āļāļąāļšāđ‚āļĄāđ€āļ”āļĨ windowing āđāļšāļš native āļ‚āļ­āļ‡ Beam āļĄāļēāļāļāļ§āđˆāļē
  • āļ›āļĢāļ°āļŠāļīāļ—āļ˜āļīāļ āļēāļž batch āļšāļ™ hardware āļ—āļĩāđˆāđ€āļ—āļĩāļĒāļšāđ€āļ—āđˆāļēāļĄāļąāļāļˆāļ°āļ”āļĩāļāļ§āđˆāļēāđƒāļ™ Spark āđ€āļ™āļ·āđˆāļ­āļ‡āļˆāļēāļāļāļēāļĢāļ›āļĢāļąāļšāđāļ•āđˆāļ‡ Catalyst/Tungsten
  • āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāđ€āļ™āđ‰āļ™āļ—āļĩāđˆ trade-offs āļĄāļēāļāļāļ§āđˆāļēāļāļēāļĢāļ›āļĢāļ°āļāļēāļĻāļ§āđˆāļē framework āđƒāļ”āđ€āļŦāļ™āļ·āļ­āļāļ§āđˆāļē
  • āļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄ hybrid āļ—āļĩāđˆāđƒāļŠāđ‰āļ—āļąāđ‰āļ‡āļŠāļ­āļ‡ framework āđ€āļ›āđ‡āļ™āđ€āļĢāļ·āđˆāļ­āļ‡āļ›āļāļ•āļīāđƒāļ™āļŠāļ āļēāļžāđāļ§āļ”āļĨāđ‰āļ­āļĄ production
  • āļāļēāļĢāđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļšāļ„āđˆāļēāđƒāļŠāđ‰āļˆāđˆāļēāļĒāļ‚āļķāđ‰āļ™āļ­āļĒāļđāđˆāļāļąāļšāļĢāļđāļ›āđāļšāļš workload: pricing per-GB āļ‚āļ­āļ‡ Dataflow vs āļāļēāļĢāļˆāļąāļ”āļŠāļĢāļĢ cluster Spark
āļŠāļēāđ€āļĨāļ™āļˆāđŒāļ›āļĢāļ°āļˆāļģāļ§āļąāļ™

āļ„āļļāļ“āļŦāļēāļšāļąāđŠāļāđƒāļ™ Data Engineering āđ€āļˆāļ­āđ„āļŦāļĄ

āđ‚āļ„āđ‰āļ”āļˆāļĢāļīāļ‡āļŦāļ™āļķāđˆāļ‡āļŠāļīāđ‰āļ™ āļšāļąāđŠāļāļ—āļĩāđˆāļ‹āđˆāļ­āļ™āļ­āļĒāļđāđˆāļŦāļ™āļķāđˆāļ‡āļˆāļļāļ” āļ§āļąāļ™āļĨāļ°āļŦāļ™āļķāđˆāļ‡āļ„āļĢāļąāđ‰āļ‡ āļĨāļ­āļ‡āđ„āļ”āđ‰āđ‚āļ”āļĒāđ„āļĄāđˆāļ•āđ‰āļ­āļ‡āļĄāļĩāļšāļąāļāļŠāļĩ

Anthony Fillion-Maillet

āđ€āļ‚āļĩāļĒāļ™āđ‚āļ”āļĒ

Anthony Fillion-Maillet

āļœāļđāđ‰āļāđˆāļ­āļ•āļąāđ‰āļ‡ SharpSkill

āđ€āļ›āđ‡āļ™āļ™āļąāļāļžāļąāļ’āļ™āļēāļŸāļđāļĨāļŠāđāļ•āļāļĄāļēāļāļ§āđˆāļē 10 āļ›āļĩ āļ”āļđāđāļĨ SharpSkill āđāļĨāļ°āļĢāļąāļšāļœāļīāļ”āļŠāļ­āļšāļ—āļļāļāļŠāļīāđˆāļ‡āļ—āļĩāđˆāđ€āļœāļĒāđāļžāļĢāđˆāļ—āļĩāđˆāļ™āļĩāđˆ

āļ­āļąāļ›āđ€āļ”āļ•āđ€āļĄāļ·āđˆāļ­ 10 āļāļąāļ™āļĒāļēāļĒāļ™ 2569

āđāļ—āđ‡āļ

#apache-beam
#spark
#data-engineering
#dataflow
#interview

āđāļŠāļĢāđŒ

āļšāļ—āļ„āļ§āļēāļĄāļ—āļĩāđˆāđ€āļāļĩāđˆāļĒāļ§āļ‚āđ‰āļ­āļ‡

dbt data transformations testing interview 2026

dbt āđƒāļ™āļ›āļĩ 2026: āļāļēāļĢāđāļ›āļĨāļ‡āļ‚āđ‰āļ­āļĄāļđāļĨ āļāļēāļĢāļ—āļ”āļŠāļ­āļš āđāļĨāļ°āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļ‡āļēāļ™

āļ„āļđāđˆāļĄāļ·āļ­ dbt āļŠāļģāļŦāļĢāļąāļšāļ§āļīāļĻāļ§āļāļĢāļ‚āđ‰āļ­āļĄāļđāļĨ: āļāļēāļĢāđāļ›āļĨāļ‡ SQL, āļāļēāļĢāļŠāļĢāđ‰āļēāļ‡āđ‚āļĄāđ€āļ”āļĨāđāļšāļšāđāļšāđˆāļ‡āļŠāļąāđ‰āļ™, āļāļĨāļĒāļļāļ—āļ˜āđŒ incremental, āļāļēāļĢāļ—āļ”āļŠāļ­āļšāļ„āļļāļ“āļ āļēāļžāļ‚āđ‰āļ­āļĄāļđāļĨ āđāļĨāļ°āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļžāļĢāđ‰āļ­āļĄāļ•āļąāļ§āļ­āļĒāđˆāļēāļ‡āđ‚āļ„āđ‰āļ”āļŠāļģāļŦāļĢāļąāļšāļ›āļĩ 2026

Apache Spark 4 new features and structured streaming

Apache Spark 4: āļŸāļĩāđ€āļˆāļ­āļĢāđŒāđƒāļŦāļĄāđˆ Structured Streaming āđāļĨāļ°āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļ‡āļēāļ™

āļŠāļģāļĢāļ§āļˆāļŸāļĩāđ€āļˆāļ­āļĢāđŒāļŠāļģāļ„āļąāļāđƒāļ™ Apache Spark 4 āļĢāļ§āļĄāļ–āļķāļ‡ ANSI SQL Mode, VARIANT data type, Real-Time Mode streaming āđāļĨāļ° transformWithState API āļžāļĢāđ‰āļ­āļĄāļ•āļąāļ§āļ­āļĒāđˆāļēāļ‡āđ‚āļ„āđ‰āļ”āđāļĨāļ°āļ„āļģāļ–āļēāļĄāļŠāļąāļĄāļ āļēāļĐāļ“āđŒāļ‡āļēāļ™āļ—āļĩāđˆāļžāļšāļšāđˆāļ­āļĒ

ETL vs ELT data pipeline architecture comparison diagram

ETL vs ELT āđƒāļ™āļ›āļĩ 2026: āļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄ Data Pipeline āļ—āļĩāđˆ Data Engineer āļ•āđ‰āļ­āļ‡āļĢāļđāđ‰

āđ€āļ›āļĢāļĩāļĒāļšāđ€āļ—āļĩāļĒāļš ETL āđāļĨāļ° ELT āļ­āļĒāđˆāļēāļ‡āļĨāļ°āđ€āļ­āļĩāļĒāļ” āļžāļĢāđ‰āļ­āļĄāļ•āļąāļ§āļ­āļĒāđˆāļēāļ‡āđ‚āļ„āđ‰āļ” dbt āđāļĨāļ° Python āđ€āļžāļ·āđˆāļ­āđ€āļĨāļ·āļ­āļāļŠāļ–āļēāļ›āļąāļ•āļĒāļāļĢāļĢāļĄ data pipeline āļ—āļĩāđˆāđ€āļŦāļĄāļēāļ°āļŠāļĄāļāļąāļšāļ—āļĩāļĄāļ‚āļ­āļ‡āļ„āļļāļ“āđƒāļ™āļ›āļĩ 2026