TGViewer
Архитектор Данных Архитектор Данных @analyticsfromzero · 1.89K subscribers
Post #396 1.45K
В Spark 4.1 появлся ... Airflow

В документации версии Spark 4.1-Preview появились так называемые Spark Declarative Pipelines (SDP)

На борту:

1️⃣ Несколько видов датасетов: Материализованные, Стриминговые, Временные
2️⃣ Пайплайн как объект. Описывается через YAML файл с SQL, Python кодом и необходимыми конфигами Спарка. Также объявляется каталог (Hive, Iceberg), с которым можно взаимодействовать и в который складывать результаты.
3️⃣ Команда spark-pipelines init с интерфейсом и аргументами как у Spark Submit. Отдельная команда spark-pipelines run.


Удобство

Пример нового кода на PySpark, который читает Kafka топик и складывает данные в таблицу в каталоге. По сути это декларативное описание (не как-сделать, а что-сделать) а-ля DAG.

from pyspark import pipelines as sdp

@sdp.table
def ingestion_st():
return (
spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "orders")
.load()
)


К объявленной таким способом таблице можно обращаться дальше по пайплайну.

На SQL и того проще

CREATE STREAMING TABLE basic_st
AS SELECT * FROM STREAM samples.nyctaxi.trips;



Или пример с несколькими синками

-- create a streaming table
CREATE STREAMING TABLE customers_us;

-- add the first append flow
CREATE FLOW append1
AS INSERT INTO customers_us
SELECT * FROM STREAM(customers_us_west);

-- add the second append flow
CREATE FLOW append2
AS INSERT INTO customers_us
SELECT * FROM STREAM(customers_us_east);


Осталось разобраться, как в этом всем провязаны семантики доставки (exactly-once, at-least-once), и куда это все полетит при смене схемы источника (Dead Letter). И понять, как устроить мониторинги и алерты работающих или сломавшихся пайплайнов.

Но ясно, что в четвертом Спарке сделать такую операцию как стриминг подхват из топиков Кафки в таблицы Айсберга будет сильно проще, чем сейчас. А то и вовсе - декларативно. Что не может не радовать.

Насладиться примерами можно в офф доке превью версии
  • 👍 15
  • 🔥 9
More from @analyticsfromzero
  1. Sep 28, 2026На инфографике вы, дорогие коллеги, можете увидеть наше светлое будущее. График составлен…
  2. Sep 28, 2026Пользовательские соглашения наконец начали читать. Раньше пользовательские соглашения никт…
  3. Sep 27, 2026Когда я даю кандидатам тестовое Даю когда человек без супер скилов и резюме. Тогда ТЗ это…
  4. Sep 26, 2026Отказ по вакансии за использование ИИ в тестовом https://t.me/bdsmmchannel/1365 Не подрыва…
  5. Sep 26, 2026SQLite написал один человек. Написал бы он в одиночку Postgres если бы использовал LLM?
  6. Sep 24, 2026Абсолютно верно Ноль возражений если воспринимать это как АйТи ПТУ (вспоминаем как это рас…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →