# Apache Beam vs Spark w 2026: Porównanie Zunifikowanych Pipeline'ów i Pytania Rekrutacyjne > Porównanie Apache Beam 2.76 i Spark 4.2 dla pipeline'ów danych. Przenośność, wydajność, pytania rekrutacyjne i wybór odpowiedniego frameworka. - Published: 2026-09-10 - Updated: 2026-09-10 - Author: Anthony Fillion-Maillet - Reading time: 5 min --- Apache Beam vs Spark to jedna z najczęstszych decyzji architektonicznych w nowoczesnej inżynierii danych. Beam 2.76 (sierpień 2026) i Spark 4.2 (lipiec 2026) obsługują zarówno przetwarzanie wsadowe, jak i strumieniowe, ale ich filozofie projektowe różnią się fundamentalnie: Beam abstrahuje silnik wykonawczy, podczas gdy Spark dostarcza ściśle zintegrowane środowisko uruchomieniowe. > **Szybka Rama Decyzyjna** > > Wybierz Beam, gdy przenośność między runnerami (Dataflow, Flink, Spark) ma znaczenie lub przy korzystaniu z Google Cloud Dataflow. Wybierz Spark przy samodzielnym zarządzaniu klastrem, korzystaniu z Databricks lub potrzebie integracji ML z MLlib. ## Model Przenośności Beam vs Zunifikowany Silnik Spark Apache Beam oddziela model programowania od wykonania. Pojedyncza definicja pipeline'u działa na [Google Cloud Dataflow](https://cloud.google.com/dataflow), Apache Flink, Apache Spark lub innych runnerach bez zmian w kodzie. Ta abstrakcja wynika z generowania przez Beam SDK przenośnej reprezentacji pipeline'u, którą interpretuje każdy kompatybilny runner. Spark 4.2 przyjmuje odwrotne podejście. DataFrame API, Structured Streaming i MLlib współdzielą ten sam optymalizator Catalyst i silnik wykonawczy Tungsten. To ścisłe powiązanie umożliwia optymalizacje takie jak [Adaptive Query Execution](https://spark.apache.org/docs/latest/sql-performance-tuning.html), które dostosowują plany w czasie wykonania na podstawie rzeczywistych statystyk danych. ```python # beam_pipeline.py import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions # Ten sam kod działa na runnerze Dataflow, Flink lub Spark options = PipelineOptions([ '--runner=DataflowRunner', # Przełącz na FlinkRunner lub 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')) ``` Odpowiednik w Spark jest bezpośrednio powiązany ze środowiskiem uruchomieniowym Spark: ```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: tryb ANSI domyślnie włączony, ściślejsze sprawdzanie typów 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() ``` Przenośność Beam umożliwia architektury niezależne od chmury, ale dodaje warstwę translacji. Bezpośrednie wykonanie Spark zwykle wykazuje niższe opóźnienia dla równoważnych operacji na tym samym sprzęcie. ## Porównanie Okienkowania i Przetwarzania Czasu Zdarzenia Oba frameworki obsługują semantykę czasu zdarzenia, ale ich API odzwierciedlają różne dziedzictwo. Model okienkowania Beam pochodzi z [Dataflow Model paper](https://research.google/pubs/pub43864/) (2015), traktując okna jako elementy pierwszej klasy pipeline'u. Spark zaadaptował swoje okienkowanie dla Structured Streaming, integrując je z DataFrame API. | Funkcja | Beam 2.76 | Spark 4.2 | |---------|-----------|----------| | Okna stałe | `FixedWindows(duration)` | `window(col, duration)` | | Okna przesuwne | `SlidingWindows(size, period)` | `window(col, size, period)` | | Okna sesji | `Sessions(gap)` | Brak natywnie (użyj `flatMapGroupsWithState`) | | Okna niestandardowe | Podklasa `WindowFn` | Ograniczone | | Obsługa spóźnionych danych | Wbudowane wyzwalacze | Opóźnienia watermarku | | Dozwolone spóźnienie | Konfiguracja per-okno | Globalny watermark | Okna sesji najwyraźniej pokazują różnicę. Beam traktuje sesje jako natywną strategię okienkowania: ```python # beam_sessions.py from apache_beam import window # 30-minutowa przerwa sesji, dopuszczenie 1 godziny spóźnionych danych windowed = ( events | 'SessionWindow' >> beam.WindowInto( window.Sessions(30 * 60), # 30 min przerwy zamyka sesję trigger=beam.trigger.AfterWatermark( early=beam.trigger.AfterProcessingTime(60), late=beam.trigger.AfterCount(1) ), allowed_lateness=3600, # Akceptuj dane do 1 godziny spóźnione accumulation_mode=beam.trigger.AccumulationMode.ACCUMULATING ) ) ``` Spark wymaga przetwarzania stanowego dla sesji: ```python # spark_sessions.py from pyspark.sql.streaming import GroupState, GroupStateTimeout def update_session(key, events, state: GroupState): # Ręczne zarządzanie sesją z API stanu 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: # Przerwa przekroczyła 30 min, emituj poprzednią sesję 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) # 30 min timeout # 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 ) ``` Dla przygotowania do rozmów kwalifikacyjnych na te tematy, warto zapoznać się z modułem [pytań rekrutacyjnych Apache Beam i Dataflow](/technologies/data-engineering/interview-questions/apache-beam-dataflow). ## Wydajność: Dane Benchmarkowe z 2026 Bezpośrednie porównania wymagają starannego przygotowania, ponieważ Beam działa na Spark jako jeden z opcjonalnych runnerów. Istotne porównanie to Beam-on-Dataflow vs natywny Spark. Najnowsze benchmarki od Databricks i Google Cloud pokazują: | Obciążenie | Spark 4.2 (Databricks) | Beam 2.76 (Dataflow) | Uwagi | |----------|------------------------|---------------------|-------| | Batch ETL (1TB Parquet) | 4.2 min | 5.1 min | Przewaga silnika Photon Spark | | Streaming (100K zdarzeń/sek) | 45ms p99 latencja | 120ms p99 latencja | Narzut autoskalowania Dataflow | | Exactly-once zapis do sink | Natywnie | Natywnie | Oba wspierają od 2024 | | Koszt (ciągłe obciążenie) | $0.12/GB przetworzone | $0.08/GB przetworzone | Ceny Dataflow Flex | Wzrost wydajności Spark 4.2 wynika z kilku funkcji: - **Tryb ANSI domyślnie**: ściślejsza semantyka SQL wyłapuje błędy wcześniej - **Typ danych VARIANT**: natywna obsługa danych półustrukturyzowanych - **Adaptive Query Execution**: optymalizacja planu w czasie wykonania - **Wsparcie Java 21**: wirtualne wątki zmniejszają narzut Beam 2.76 odpowiada: - **Wsparcie runnera Flink 2.0**: produkcyjne przetwarzanie stanowe streaming - **Persystencja offsetów CDC**: DebeziumIO FileSystemOffsetRetainer - **Integracja ADK**: wsparcie Google Agent Development Kit w Python SDK ## Popularne Pytania Rekrutacyjne: Beam vs Spark Rozmowy techniczne na stanowiska inżyniera danych często porównują te frameworki. Poniżej pytania pojawiające się na rozmowach w 2026 roku wraz z oczekiwaną głębią odpowiedzi. **P1: Kiedy wybrać Beam zamiast natywnego Spark?** Oczekiwane punkty odpowiedzi: - Wdrożenia multi-cloud lub hybrydowe, gdzie kod pipeline'u musi działać bez zmian - Środowiska Google Cloud, gdzie Dataflow zapewnia zarządzaną infrastrukturę - Złożona semantyka czasu zdarzenia wymagająca okien sesji lub niestandardowych wyzwalaczy - Zespoły z istniejącym doświadczeniem Beam z tła Dataflow Odpowiedź-czerwona flaga: "Beam jest zawsze lepszy, bo jest przenośny." Przenośność ma swoje koszty. **P2: Jak abstrakcja runnera Beam wpływa na debugowanie?** Oczekiwane punkty odpowiedzi: - Stack trace'y odwołują się zarówno do Beam SDK, jak i implementacji runnera - Optymalizacje specyficzne dla runnera mogą zachowywać się inaczej (checkpointy Flink vs checkpointy Spark) - Metrics API zapewnia zunifikowany monitoring, ale dashboardy runnerów pokazują różne szczegóły - Testowanie z DirectRunner przed wdrożeniem na produkcyjny runner **P3: Wyjaśnij semantykę exactly-once w obu frameworkach.** Oczekiwana odpowiedź: ```python # Beam: exactly-once poprzez gwarancje runnera # Dataflow zapewnia exactly-once zarówno dla źródeł jak i sinków # SDK obsługuje deduplikację i koordynację checkpointów with beam.Pipeline() as p: (p | beam.io.ReadFromPubSub(subscription='...') # Exactly-once odczyt | beam.Map(process) | beam.io.WriteToBigQuery(...) # Exactly-once zapis z retries ) # Spark: exactly-once poprzez checkpointing i idempotentne sinki spark.readStream \ .format("kafka") \ .load() \ .writeStream \ .option("checkpointLocation", "/checkpoint") # Odzyskiwanie stanu .foreachBatch(idempotent_write) # Deduplikacja na poziomie aplikacji .start() ``` Więcej tematów dotyczących streamingu znajdziesz w module [PySpark](/technologies/data-engineering/interview-questions/pyspark). **P4: Jak przeprowadzić migrację zadania wsadowego Spark do Beam?** To pytanie testuje zrozumienie obu API. Kluczowe punkty: 1. Mapowanie operacji DataFrame na PCollections i transformacje 2. Zastąpienie `spark.read` odpowiednimi konektorami I/O Beam 3. Konwersja UDF na funkcje `beam.Map` lub `beam.ParDo` 4. Inna obsługa partycjonowania (`Reshuffle` Beam vs `repartition` Spark) 5. Testowanie z DirectRunner przed wdrożeniem na produkcyjny runner ## Wybór Odpowiedniego Narzędzia: Matryca Decyzyjna Wybór zależy bardziej od kontekstu organizacyjnego niż możliwości technicznych. Oba frameworki kompetentnie obsługują większość zadań inżynierii danych. | Czynnik | Faworyzuje Beam | Faworyzuje Spark | |--------|-------------|-------------| | Dostawca chmury | Google Cloud | AWS EMR, Databricks, on-prem | | Doświadczenie zespołu | Istniejące doświadczenie Dataflow | Istniejące umiejętności Spark/PySpark | | Integracja ML | Ograniczona (oddzielne narzędzia) | MLlib, Spark ML | | Analiza interaktywna | Nie zaprojektowany do tego | Spark SQL, notebooki | | Okna sesji | Natywne wsparcie | Ręczne zarządzanie stanem | | Model kosztowy | Pay-per-use (Dataflow) | Provisioning klastra | | Vendor lock-in | Niższy (wielu runnerów) | Wyższy (kod specyficzny dla Spark) | Dla [decyzji o wzorcach ETL/ELT](/technologies/data-engineering/interview-questions/etl-elt-patterns), należy rozważyć wolumen danych i wymagania dotyczące opóźnień przed wyborem frameworka. ## Rzeczywista Architektura: Podejście Hybrydowe Wiele organizacji używa obu frameworków. Popularny wzorzec: ``` [Streaming Ingestion] [Batch Processing] [ML Training] | | | Beam/Dataflow Spark on Databricks Spark MLlib | | | v v v BigQuery <----- dbt -----> Delta Lake -----> Model Registry ``` Beam obsługuje ingestion streamingowy, gdzie autoskalowanie Dataflow dopasowuje się do wzorców ruchu. Spark przetwarza obciążenia wsadowe, gdzie ekonomika klastra faworyzuje ciągłe obliczenia. Oba zasilają zunifikowany data warehouse. Ta architektura pojawia się w [tutorialu Apache Spark](/blog/data-engineering/apache-spark-pyspark-data-pipelines-tutorial) ze szczegółami implementacji. ## Kluczowe Wnioski dla Wyboru Beam vs Spark - Beam 2.76 zapewnia przenośność runnera między Dataflow, Flink 2.0 i Spark, wymieniając część wydajności na elastyczność wdrożenia - Spark 4.2 dostarcza ściślejszą integrację między SQL, streamingiem i obciążeniami ML z funkcjami takimi jak typy VARIANT i Adaptive Query Execution - Okna sesji i złożone wyzwalacze faworyzują natywny model okienkowania Beam - Wydajność wsadowa na równoważnym sprzęcie zwykle faworyzuje Spark dzięki optymalizacji Catalyst/Tungsten - Pytania rekrutacyjne skupiają się na kompromisach, a nie na deklarowaniu jednego frameworka jako lepszego - Architektury hybrydowe używające obu frameworków są powszechne w środowiskach produkcyjnych - Porównanie kosztów zależy od wzorców obciążenia: ceny per-GB Dataflow vs provisioning klastra Spark --- Source: SharpSkill (https://sharpskill.dev), tech interview preparation for your real stack. HTML version of this page: https://sharpskill.dev/pl/blog/data-engineering/apache-beam-vs-spark-2026-unified-pipelines-interview-questions