惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

宝玉的分享
宝玉的分享
AWS News Blog
AWS News Blog
Y
Y Combinator Blog
云风的 BLOG
云风的 BLOG
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
F
Full Disclosure
H
Help Net Security
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
A
About on SuperTechFans
J
Java Code Geeks
Jina AI
Jina AI
GbyAI
GbyAI
酷 壳 – CoolShell
酷 壳 – CoolShell
爱范儿
爱范儿
美团技术团队
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
Latest news
Latest news
Vercel News
Vercel News
博客园 - 【当耐特】
P
Privacy & Cybersecurity Law Blog
P
Proofpoint News Feed
阮一峰的网络日志
阮一峰的网络日志
V
Vulnerabilities – Threatpost
Stack Overflow Blog
Stack Overflow Blog
Hugging Face - Blog
Hugging Face - Blog
D
Docker
Microsoft Security Blog
Microsoft Security Blog
博客园_首页
S
Securelist
WordPress大学
WordPress大学
S
Secure Thoughts
博客园 - 聂微东
Cloudbric
Cloudbric
Help Net Security
Help Net Security
腾讯CDC
T
Threat Research - Cisco Blogs
T
Tor Project blog
L
LINUX DO - 热门话题
Project Zero
Project Zero
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
博客园 - Franky
N
Netflix TechBlog - Medium
小众软件
小众软件
Cyberwarzone
Cyberwarzone
量子位
MyScale Blog
MyScale Blog
W
WeLiveSecurity
MongoDB | Blog
MongoDB | Blog
I
InfoQ
M
MIT News - Artificial intelligence

Все публикации подряд на Хабре

Ловим музу за клавиатуру: как айтишнику стать автором Что умеет Midjourney в 2026? Мой немного грустный разбор этого шикарного инструмента Никто не любит писать тесты, но ИИ может исправить это IPv8 выглядит как мечта. Поэтому почти наверняка не взлетит Производители вернули в продажу материнки с DDR3. Что происходит? Управление агентом с телефона через Telegram теперь в KodaCode От координации к лидерству: как меняется роль руководителя разработки Я сделала родителям бизнес вместо пенсии: зарабатываем 70 тысяч, мама не даёт продать В три раза быстрее приемка товара и оптимизация трудозатрат на 73%: как «РСТ-Инвент» помог Gulliver Group ИИ-шечный мир победил? О влиянии искусственного интеллекта на игропром T-TOPS: Как распутать гордиев узел проекта после выхода в прод (меч не понадобится) Кремль снижает давление на Телеграмм пока Европа строит интернет по паспорту Как CEO, CTO и CIO за 8 часов собрали ИИ-директора, который умеет держать позицию под давлением Как (не) потерять домен за выходные Вместо 8 разных VPS: как я организовал практику студентам на одном сервере Почему твой Open Source проект не замечают? R&D: искусство управления неопределенностью в разработке AI-дефляция: вакансий для разработчиков больше, а рост зарплат — худший за 15 лет Мы отдали управление роботами OpenClaw. Что из этого вышло Галактический ID: система идентификации для всех форм разумной жизни Шесть основ бизнес-анализа: начинаем с вопроса «Кто в игре?» Код-ревью, в котором дело не в коде Данные переехали. Команда — нет Системной подход к сдаче OSWE в 2025 Почему комната управления реактором покрашена в цвет морской пены LLM-агент для поиска свободных доменов: автоматизируем подбор Когда, зачем и как правильно начинать новую сессию в Claude Code? Как я заставил нейросеть писать макросы для FreeCAD Анатомия ИИ‑агента для подбора персонала. От тысячи резюме к топ‑10 за минуты Опыт разработчика как экономика внимания Автономность как точка невозврата: кто будет субъектом в цифровом будущем Обучение ИИ в «диких» условиях: как рутинные действия превращаются в датасеты Как измерить LLM для задач кибербеза: обзор открытых бенчмарков Где хранить код? Сравнение GitHub, GitLab и Bitbucket Математика объясняет, почему нормальное распределение встречается повсюду Почему ваш FinOps не работает: 12 тезисов от практиков Как подписать проектную документацию УКЭП с использованием бесплатных лицензий Pilot Адаптивное администрирование Sigla Vision Я грузил уран в бочки, а потом 20 лет строил ИТ в атомной отрасли Чем позвонить с Эвереста? История и обзор спутниковой связи. Часть 2 Как языковая модель помогает контролировать качество инструктажей по охране труда в металлургии Как не передать на desktop свой IP в РКН Анатомия SAP Privileges: как устроено управление правами в macOS MoneyDev: Сказка про три главных слова Обновлённый токенизатор видео K-VAE 2.0 от Сбера Как сделать диспетчеризацию дома на 1284 квартиры почти бесплатно Как мы разогнали железную дорогу Мы дали агентам рутину. Теперь надо решить — что делать с освободившимся временем Токсичный контент, промпт-хакинг и защита ИИ — всё о Guardrails для LLM Умный город начинается с точного взгляда: как «Фалькон Тех» меняет пространство к лучшему Навайбкодил приложение для анализа графов Почему Дюну так интересно читать? Упрощаем работу с рутиной или как стать Гендальфом Белым Деконструкция Go: CPU, RAM и что там происходит. Go Assembler база. Часть 1.1 Какие профессии исчезнут из-за ИИ, а какие появятся? И что с этим делать Как мы построили IT-отдел, где хочется расти: архитектурные встречи, прозрачные метрики и книжные подарки Rufler: Делаем из Claude Code автономный рой через один YAML-конфиг Sing-box и белый список приложений Как построить надёжный обмен сообщениями в микросервисах: лучшие практики для enterprise OpenAI строит MLM-пирамиду, а McKinsey и Accenture помогают ей в этом Дом, который не построил Фишер (Часть 2) «Сверхзвуковой математик» против «Вдумчивого логиста»: битва алгоритмов 3D-упаковки Мультимодальные модели – грубый и дорогой инструмент Разговоры ничего не стоят. Код тоже Проверки физических лиц: с кого начнет ФНС Топ-10 бесплатных нейросетей для создания видео в 2026 году Первые слои кода: как наши решения сегодня определяют архитектуру ИИ на десятилетия Разработка нового статического анализатора: PVS-Studio JavaScript Поиск уязвимостей ПО: базовый минимум или роскошный максимум Почему оценка персонала не работает как инструмент управления Как мы разработали ИИ-ассистента и сократили рутину продуктовой команды на 50% Как я ушел из найма, нажарил косточек и продал на маркетплейсах на 168 млн в год Когда 1С:ERP уже внедрена, а нормального производственного плана всё ещё нет Как я сделал Claude мультимодальным, подключив к нему Qwen Omni Как приглашение на вакансию мечты превращается в атаку Infrastructure as Code: философия и лучшие практики IaC Тестируем Yandex Code Assistant на задаче, в которой нужно хранить секреты nxs-universal-chart v3.0: новое поколение универсального Helm-чарта Callback Injection: Техника, которая отправила Microsoft Defender в глухой нокаут «Все идеи на стол»: митап как способ вывести проект из тупика Сегодня я узнал нечто новое о GPU благодаря багу в своей игре Как заставить LLM ̶ ̶г̶а̶л̶л̶ю̶ ̶ эволюционировать Карта событий как фундамент аналитики: практический кейс для E-commerce Что выбрать для AI: x86, ARM или RISC-V? Дайджест железа за март Роль соматических мутаций в развитии аутоиммунных заболеваний: путь к избирательной терапии Mythos от Anthropic — тревожный сигнал для всех, а не только для банков Guardrails для LLM на Java: как приручить промпт‑инъекции и токсичные ответы Green-VLA: как мы собрали VLA-модель для реального антропоморфного робота и не потеряли обобщение Финансовая гонка вооружений: почему умные люди добровольно в ней участвуют Эра ИИ-агентов наступила: выбираем лучшего цифрового сотрудника # Практический опыт внедрения WinCC Redundancy на производственном предприятии Сделал MVP за 3 дня, а потом неделю прикручивал оплату. Оно того стоило? Физика против Маска: почему Starship V3 может оказаться ещё одной катастрофой Нефть Венесуэлы: крупнейшие запасы в мире, но не крупнейшая нефтяная держава JPA 4. Переосмысление Hibernate Почему зеркальная фотокамера Nikon D5 десятилетней давности идеально подошла для миссии «Артемида-2» Проект «Уровень-Спутник» или как мы сделали платформу для гидрологов «Замедлиться, чтобы ускориться»: почему ИИ повышает цену ошибок в требованиях и архитектуре Как с нуля поднять трафик IT-компании на 1657% при бюджете 55 тыс. и выжить Pixel-perfect Downsampling — идеальная отрисовка 50 миллионов точек без потерь
4 YAML-файла вместо PySpark: как аналитикам строить пайплайны без разработчиков
Kiril Kazlou · 2026-04-16 · via Все публикации подряд на Хабре

Средний

10 мин

5.1K

Всем привет! На связи Кирилл Козлов, data‑инженер Mindbox. Наша команда регулярно пересчитывает бизнес‑метрики для клиентов. Для этого нам приходится формировать витрины данных для биллинга и аналитики на основе десятков источников. 

Долгое время мы обрабатывали данные для расчетов на PySpark — инструменте, с которым сложно работать без опыта программирования на Python. Чтобы создать любой пайплайн, приходилось привлекать разработчиков. Это затягивало процесс на несколько недельных спринтов.

В статье расскажу, как мы построили внутреннюю data‑платформу, где аналитик или продакт может создать регулярно обновляемый пайплайн, описав его в четырех YAML‑файлах.

Почему PySpark тормозил создание пайплайнов

Опишу нашу боль на примере типичного пайплайна — элементарного подсчета MAU.

На первый взгляд, задачу можно было бы решить с помощью простого SQL‑запроса — COUNT(DISTINCT customerId) по нескольким таблицам за период. Но из‑за инфраструктурного обвеса — PySpark, Airflow DAG, Spark‑ресурсы, тесты — пришлось отдать задачу разработчикам. В итоге мы ждали пайплайн подсчета MAU целую неделю.

На подготовку расчета каждой новой метрики уходило от 1 до 3 недель. И каждый раз приходилось повторять один и тот же сценарий: 

  1. Аналитик формулировал бизнес‑требования, искал свободного разработчика, передавал ему контекст задачи. 

  2. Разработчик уточнял детали, писал код на Python, проходил ревью, настраивал DAG и деплоил результат.

Мы же хотели ускорить этот процесс и решать такие задачи силами аналитиков и менеджеров продукта, которые лучше всех понимают бизнес‑логику, свободно владеют SQL и YAML, но не работают с Python и PySpark.

Раньше так создавали пайплайны на PySpark

Раньше так создавали пайплайны на PySpark

Чем заменили PySpark: Python не нужен, достаточно YAML и SQL

Чтобы использовать декларативный подход для создания пайплайнов, мы разделили весь data‑слой на три части и для каждой подобрали инструмент:

  1. dlt (data load tool) загружает данные из внешних API и БД в объектное хранилище. Конфигурируется YAML‑файлом без кода.

  2. dbt (data build tool) on Trino преобразует данные с помощью SQL. Связывает модели через ref(), сам строит граф зависимостей и управляет инкрементальными обновлениями.

  3. Airflow + Cosmos оркестрирует пайплайны. DAG для Airflow генерируется автоматически из dag.yaml и dbt‑проекта.

Trino мы уже использовали как query engine для ad‑hoc‑запросов и подключали его к Superset для BI. Он уже тогда показал себя хорошо: для запросов с базовой логикой обрабатывал огромные объемы данных быстрее и с меньшим потреблением ресурсов, чем Spark. Кроме того, Trino из коробки поддерживает федеративный доступ к нескольким хранилищам из одного SQL‑запроса. Для 90% наших пайплайнов Trino подходил отлично.

Взаимодействие новых инструментов с потоками данных в нашей data-платформе

Взаимодействие новых инструментов с потоками данных в нашей data-платформе

Как загружаем данные: dlt.yaml

Первый yaml‑файл описывает, откуда и как загружать данные для дальнейшей обработки. Вот реальный пример того, как загружаются данные биллинга из внутреннего API:

product: sg-team
feature: billing
schema: billing_tarification
dag:
  dag_id: dlt_billing_tarification
  schedule: "0 4 * * *"
  description: "Daily refresh of tarification data"
  tags:
    - billing

alerts:
  enabled: true
  severity: warning

source:
  type: rest_api
  client:
    base_url: "https://internal-api.example.com"
    auth:
      type: bearer
      token: dlt-billing.token

  resources:
    - name: tarification_data
      endpoint:
        path: /tarificationData
        method: POST
        json:
          firstPeriod: "{{ previous_month_date }}"
          lastPeriod: "{{ previous_month_date }}"
          pricingPlanLine: CurrentPlan
      write_disposition: replace
      processing_steps:
        - map: dlt_custom.billing_tarification_data.map

    - name: charges_raw
      columns:
        staffUserName:
          data_type: text
          nullable: true
      endpoint:
        path: /data-feed/charges
        method: POST
        json:
          firstPeriod: "{{ previous_month_date }}"
          lastPeriod: "{{ previous_month_date }}"
      write_disposition: replace

    - name: discounts_raw
      endpoint:
        path: /data-feed/discounts
        method: POST
        json:
          firstPeriod: "{{ previous_month_date }}"
          lastPeriod: "{{ previous_month_date }}"
      write_disposition: replace

В файле описаны четыре ресурса из одного API. Для каждого указывается endpoint, параметры запроса и стратегии записи: в нашем случае replace означает «перезаписывать каждый раз». В этом конфиге можно добавлять шаги обработки processing_steps, описывать типы колонок и настраивать алерты.

В результате у нас всего 40 строк YAML‑конфига. Без dlt каждый такой коннектор — это Python‑скрипт с запросами, обработкой пагинации, ретраями, сериализацией в Delta Table формат и загрузкой в хранилище.

Как трансформируем данные с помощью SQL: dbt_project.yaml и sources.yaml

Следующим шагом конфигурируем dbt‑модель. Для Trino это SQL‑запросы.

Приведу пример конфигурации для расчета MAU. Так выглядит подготовка событий из одного источника:

-- int_mau_events_visits.sql
{{ config(materialized='table') }}

WITH period AS (
    SELECT
        YEAR(CURRENT_DATE - INTERVAL '5' MONTH) AS start_year,
        MONTH(CURRENT_DATE - INTERVAL '5' MONTH) AS start_month,
        YEAR(CURRENT_DATE) AS end_year,
        MONTH(CURRENT_DATE) AS end_month
),

merged_customers AS (
    SELECT * FROM {{ ref('int_merged_customers') }}
),

actual_customers AS (
    SELECT * FROM {{ ref('int_actual_customers') }}
),

events AS (
    SELECT
        src._tenant,
        src.unmergedCustomerId,
        'visits' AS src_type,
        src.endpoint
    FROM {{ source('final', 'customerstracking_visits') }} src
    CROSS JOIN period p
    WHERE src.unmergedCustomerId IS NOT NULL
        AND ((src.timestamp_year > p.start_year)
             OR (src.timestamp_year = p.start_year
                 AND src.timestamp_month >= p.start_month))
        AND ((src.timestamp_year < p.end_year)
             OR (src.timestamp_year = p.end_year
                 AND src.timestamp_month <= p.end_month))
),

events_with_customer AS (
    SELECT
        e._tenant,
        COALESCE(mc.mergedCustomerId, e.unmergedCustomerId) AS customerId,
        e.src_type,
        e.endpoint
    FROM events e
    LEFT JOIN merged_customers mc
        ON e._tenant = mc._tenant
        AND e.unmergedCustomerId = mc.unmergedCustomerId
)

SELECT
    ewc._tenant, ewc.customerId, ewc.src_type, ewc.endpoint
FROM events_with_customer ewc
WHERE EXISTS (
    SELECT 1 FROM actual_customers ac
    WHERE ewc._tenant = ac._tenant AND ewc.customerId = ac.customerId
)

Все 10 источников событий построены по одному и тому же принципу. Меняется только таблица‑источник и фильтры. Далее модели собираются в единый поток:

-- int_mau_events.sql (объединение всех источников)
SELECT * FROM {{ ref('int_mau_events_inapps_targetings') }}
UNION ALL
SELECT * FROM {{ ref('int_mau_events_inapps_clicks') }}
UNION ALL
SELECT * FROM {{ ref('int_mau_events_visits') }}
UNION ALL
SELECT * FROM {{ ref('int_mau_events_orders') }}
-- ...и ещё 6 источников

И наконец финальная витрина, где все считается:

-- mau_period_datamart.sql
{{ config(
    materialized='incremental',
    incremental_strategy='merge',
    unique_key=['_tenant', 'start_year', 'start_month', 'end_year', 'end_month']
) }}

{%- set months_back = var('months_back', 5) | int -%}

WITH period AS (
    SELECT
        YEAR(CURRENT_DATE - INTERVAL '{{ months_back }}' MONTH) AS start_year,
        MONTH(CURRENT_DATE - INTERVAL '{{ months_back }}' MONTH) AS start_month,
        YEAR(CURRENT_DATE) AS end_year,
        MONTH(CURRENT_DATE) AS end_month
),

events_resolved AS (
    SELECT * FROM {{ ref('int_mau_events') }}
),

metrics_by_tenant AS (
    SELECT
        er._tenant,
        COUNT(DISTINCT CASE WHEN src_type = 'visits'
              THEN customerId END) AS CustomersTracking_Visits,
        COUNT(DISTINCT CASE WHEN src_type = 'orders'
              THEN customerId END) AS ProcessingOrders_Orders,
        COUNT(DISTINCT CASE WHEN src_type = 'mailings'
              THEN customerId END) AS Mailings_MessageStatuses,
        -- ...остальные метрики
        COUNT(DISTINCT customerId) AS MAU
    FROM events_resolved er
    GROUP BY er._tenant
)

SELECT m.*, p.start_year, p.start_month, p.end_year, p.end_month
FROM metrics_by_tenant m
CROSS JOIN period p

Для конфигурации витрины используем incremental_strategy='merge'. В этом случае dbt сам генерирует merge-запрос, подставляя unique_key для upsert. Не нужно вручную реализовывать инкрементальную загрузку.

Чтобы связать модели в один проект, собираем конфиг dbt_project.yaml:

name: mau_period
version: '1.0.0'
models:
  mau_period:
    +on_table_exists: replace
    +on_schema_change: append_new_columns

И файл sources.yaml, который описывает входные таблицы:

sources:
  - name: final
    database: data_platform
    schema: final
    tables:
      - name: inapps_targetings_v2
      - name: inapps_clicks_v2
      - name: customerstracking_visits
      - name: processingorders_orders
      - name: cdp_mergedcustomers_v2
      # ...

В итоге получаем ту же бизнес-логику, что и в PySpark, но на чистом SQL:

sources.yaml вместо typedspark-схем,

{{ ref() }} и {{ source() }} вместо .get_table(),

— автоматический порядок выполнения по графу зависимостей вместо ручного подбора Spark-ресурсов.

Как конфигурируем Airflow: dag.yaml

Четвертый файл конфигурации определяет, когда и как Airflow запустит пайплайн:

product: sg-team
feature: billing
schema: mau
schedule: "15 21 * * *"  # каждый день в 00:15 MSK

params:
  - name: start_date
    description: "Start date (YYYY-MM-DD). Leave empty for auto"
    default: ""
  - name: end_date
    description: "End date (YYYY-MM-DD). Leave empty for auto"
    default: ""
  - name: months_back
    description: "Months to look back (default: 5)"
    default: 5

alerts:
  enabled: true
  severity: warning

Затем наш Python‑скрипт парсит dag.yaml и dbt_project.yaml и через библиотеку Cosmos генерирует полноценный Airflow DAG. Это единственный кусок кода на Python. Он пишется один раз и работает для всех dbt‑проектов. Вот ключевая часть этого скрипта:

def _build_dbt_project_dags(project_path: Path, environ: dict) -> list[DbtDag]:
    config_dict = yaml.safe_load(dag_config_path.read_text())
    config = DagConfig.model_validate(config_dict)

    # Параметры из YAML → Airflow Params
    params = {}
    operator_vars = {}
    for param in config.params:
        params[param.name] = Param(
            default=param.default if param.default is not None else "",
            description=param.description,
        )
        operator_vars[param.name] = f"{{{{ params.{param.name} }}}}"

    # Cosmos создаёт DAG из dbt-проекта
    with DbtDag(
        dag_id=f"dbt_{project_path.name}",
        schedule=config.schedule,
        params=params,
        project_config=ProjectConfig(dbt_project_path=project_path),
        profile_config=ProfileConfig(
            profile_name="default",
            target_name=project_name,
            profile_mapping=TrinoLDAPProfileMapping(
                conn_id="trino_default",
                profile_args={
                    "database": profile_database,
                    "schema": profile_schema,
                },
            ),
        ),
        operator_args={"vars": operator_vars},
    ) as dag:
        # Создаём схему перед запуском моделей
        create_schema = SQLExecuteQueryOperator(
            task_id="create_schema",
            conn_id="trino_default",
            sql=f"CREATE SCHEMA IF NOT EXISTS {profile_database}.{profile_schema} ...",
        )
        # Привязываем к корневым задачам
        for unique_id, _ in dag.dbt_graph.filtered_nodes.items():
            task = dag.tasks_map[unique_id]
            if not task.upstream_task_ids:
                create_schema >> task

Cosmos читает manifest.json из dbt-проекта, разбирает граф зависимостей моделей и создаёт отдельную Airflow-задачу для каждой модели. Зависимости между задачами выстраиваются автоматически на основе ref() в SQL.

Как создаем пайплайны, не привлекая разработчиков

Теперь, когда аналитик хочет создать новый регулярный пайплайн, он может собрать его за несколько шагов:

Шаг 1. Создает папку в репозитории dbt-projects/my_new_pipeline/.

Шаг 2. Если нужна загрузка из внешнего источника, пишет YAML-конфиг для dlt. 

Шаг 3. Формирует SQL-модели в папке models/ и описывает источники в sources.yaml.

Шаг 4. Создаёт dbt_project.yaml и dag.yaml.

Шаг 5. Пушит в Git, проходит ревью, мержит.

CI/CD собирает dbt‑проект и отправляет артефакты на S3. Airflow читает DAG‑файлы оттуда, Cosmos парсит dbt‑проект и генерирует граф задач. По расписанию dbt запускает модели на Trino в правильном порядке. В итоге имеем обновленную витрину в хранилище, доступную в Superset.

Что изменилось после перехода на новую архитектуру пайплайнов

PySpark

dlt, Trino, Airflow+Cosmos

Пайплайн создавался 1–3 недели

Аналитик создает пайплайн за 1 день

Чтобы создать пайплайн, приходилось обращаться к разработчикам

Аналитики решают большую часть задач самостоятельно

Пайплайны требовали большой кодовой базы на Python

Пайплайны описывают в SQL и YAML, а объем кода уменьшился в 2 раза

Чтобы аналитик мог создавать пайплайны своими силами, он должен понимать концепции ref() и source(), разницу между table и incremental и знать основы Git. Для этого мы провели несколько внутренних воркшопов и составили пошаговые инструкции для каждого типа задач.

Как создаем пайплайны на dbt + Trino

Как создаем пайплайны на dbt + Trino

Почему новые инструменты не заменяют PySpark полностью

Для 10% наших пайплайнов PySpark все еще остается единственным вариантом, если трансформация не ложится в SQL. dbt поддерживает Jinja‑макросы, но это не замена полноценному Python. Кроме того, было бы нечестно умалчивать об ограничениях новых инструментов.

dlt + Delta: экспериментальный upsert. Мы используем Delta-формат в хранилище. У dlt коннектор для Delta помечен как experimental, поэтому merge-стратегия не работала. Пришлось придумать обходные пути: в одних случаях использовали replace вместо merge и теряли инкрементальность, в других — писали кастомные processing_steps.

Низкая отказоустойчивость Trino. У Trino есть механизм fault tolerance, но он работает через запись промежуточных результатов на S3. При наших терабайтных объемах данных использовать механизм нецелесообразно — дорого из-за огромного количества S3-операций. Если не включать fault tolerance, то при падении воркера Trino роняет весь запрос. В отличие от Spark, который перезапускает только упавшую задачу. Мы решили эту проблему через ретраи в DAG и декомпозицию тяжелых моделей на цепочку промежуточных.

UDF и кастомная логика. В Spark кастомная логика пишется на Python прямо в пайплайне, это удобно. С новой архитектурой делать такое сложнее. dbt поверх Trino не спасает: Jinja генерирует лишь SQL, а Python-модели работают только со Snowflake, Databricks и BigQuery. В самом Trino можно писать UDF, но только на Java со всеми вытекающими последствиями: отдельный репозиторий, сборка, деплой JAR на все воркеры. Таким образом, когда трансформация не ложится в SQL, UDF для Trino становится либо неподдерживаемым SQL-монстром, либо отдельным скриптом с потерей lineage.

Что планируем улучшить: тесты, шаблоны для моделей и обучение

Доделать тестирование. В PySpark у нас было отлажено тестирование пайплайнов, а в новой архитектуре тестов пока не хватает. В свежих версиях dbt появилась возможность юнит-тестирования. Теперь можно проверять логику SQL-моделей на мок-данных, не поднимая полный пайплайн. Хотим добавить dbt-тесты и на уровне моделей, и как отдельный слой мониторинга.

Внедрить шаблоны для типовых задач. Многие наши dbt-модели похожи. Один конфиг может описывать десяток моделей с одинаковым паттерном: разница только в таблице-источнике и применяемых фильтрах. Планируем вынести общую логику в dbt-макросы.

Расширить круг пользователей платформы. Мы хотим, чтобы больше инженеров и аналитиков могли самостоятельно работать с данными. Планируем регулярно проводить внутреннее обучение, готовить документацию и онбординг-гайды, чтобы новые пользователи могли быстро подключаться к платформе и строить свои модели.