Решение проблем Apache DataFusion
Частые ошибки и проблемы при работе с Apache DataFusion — симптомы, причины и пошаговые решения.
Область
Категория
Симптомы
- RecordBatch не создаётся — ошибка несовпадения типов колонок
- При построении RecordBatch из массивов Arrow возвращается ArrowError
- Тип данных в ArrayRef не соответствует ожидаемому в Schema
Причина
Схема (Schema), переданная в RecordBatch::try_new(), объявляет один тип для поля, а реальный ArrayRef содержит другой. Часто возникает при ручном создании батчей или в TableProvider, когда schema() возвращает одну схему, а scan() — батчи с другой.
Решение
- Сравните Schema из
schema()с фактическими типами ArrayRef в батче - Используйте
arrow::compute::cast()для приведения типов, если данные совместимы - Убедитесь, что
TableProvider::schema()иTableProvider::scan()возвращают одинаковую схему
Связанные уроки:
Симптомы
- Запрос прерывается с ошибкой о превышении лимита памяти
- MemoryPool возвращает ResourcesExhausted при аллокации
- Крупные JOIN или агрегации падают, хотя данные помещаются в RAM
Причина
MemoryPool в DataFusion ограничивает общий объём памяти для запроса или сессии. При обработке больших объёмов данных (JOIN, сортировка, агрегация) промежуточные буферы превышают лимит пула.
Решение
- Увеличьте лимит:
ctx.runtime_env().memory_pool.set_limit(bytes) - Используйте
FairSpillPoolвместоGreedyMemoryPoolдля более равномерного распределения - Уменьшите
batch_sizeв RuntimeConfig для снижения пикового потребления - Добавьте repartition перед JOIN для лучшего распределения данных
Связанные уроки:
Симптомы
- SQL-запрос с CAST() или неявным приведением типов падает
- Ошибка при попытке прочитать строковую колонку как числовую
- Parquet/CSV файл содержит невалидные значения для целевого типа
Причина
Arrow выполняет строгое приведение типов. Если строковое значение не может быть распарсено в целевой тип (например, 'abc' → Int32), cast завершается ошибкой вместо возврата NULL.
Решение
- Используйте
TRY_CAST(col AS INT)вместоCAST()— возвращает NULL при ошибке конвертации - Предварительно очистите данные:
SELECT * FROM t WHERE col ~ '^[0-9]+$' - При чтении CSV задайте правильный schema в
CsvReadOptions::new().schema(&schema)
Связанные уроки:
Симптомы
- SQL-запрос к зарегистрированной таблице возвращает 'table not found'
- После перезапуска программы таблицы пропадают из каталога
- ctx.sql("SELECT * FROM my_table") падает при первом вызове
Причина
Таблица не зарегистрирована в текущем SessionContext. В DataFusion регистрация таблиц существует только в памяти — при каждом новом SessionContext нужно регистрировать заново. Также возможна ошибка в имени каталога/схемы.
Решение
- Зарегистрируйте таблицу до выполнения запроса:
ctx.register_parquet("my_table", "path.parquet", ...).await? - Проверьте доступные таблицы:
ctx.sql("SHOW TABLES").await? - Укажите полное имя:
SELECT * FROM datafusion.public.my_table - Для персистентного каталога реализуйте свой CatalogProvider
Связанные уроки:
Симптомы
- JOIN-запрос падает с ошибкой 'Ambiguous reference'
- Колонка присутствует в нескольких таблицах, и DataFusion не может выбрать
- Ошибка при SELECT без полного квалификатора таблицы
Причина
Несколько таблиц в JOIN содержат колонку с одинаковым именем. DataFusion требует явного указания таблицы-источника для неоднозначных ссылок, как и стандартный SQL.
Решение
- Используйте полную квалификацию:
SELECT a.id, b.id FROM table_a a JOIN table_b b ON a.id = b.id - Задайте алиасы таблицам:
FROM orders AS o JOIN products AS p - В DataFrame API:
left.join(right, JoinType::Inner, &["id"], &["id"], None)?
Связанные уроки:
Симптомы
- Оконная функция с RANGE BETWEEN возвращает 'not implemented'
- Работает с ROWS BETWEEN, но падает с RANGE BETWEEN
- Сложные оконные выражения не поддерживаются
Причина
DataFusion поддерживает RANGE BETWEEN только для ограниченного набора типов (numeric, temporal). Для строковых и других типов рамка RANGE не реализована. Это известное ограничение, а не баг.
Решение
- Замените
RANGE BETWEENнаROWS BETWEEN— в большинстве случаев результат эквивалентен - Для temporal-колонок убедитесь, что тип — Timestamp, а не Utf8
- Проверьте текущий список поддерживаемых функций:
SELECT * FROM information_schema.df_settings
Связанные уроки:
Симптомы
- EXPLAIN ANALYZE показывает 0 строк на одном из шагов плана
- Запрос возвращает пустой результат, хотя данные есть
- FilterExec отсекает все строки — предикат не совпадает с данными
Причина
Фильтр не совпадает с данными из-за несовпадения типов (например, строковый '1' != числовой 1), регистра (case-sensitive сравнение), или предикат применяется к пустому батчу после предыдущей операции.
Решение
- Запустите
EXPLAIN ANALYZEи найдите шаг с output_rows=0 - Проверьте типы колонок:
DESCRIBE my_table - Убедитесь, что литерал имеет правильный тип:
WHERE id = CAST('1' AS INT)вместоWHERE id = '1' - Используйте
EXPLAIN(без ANALYZE) для просмотра логического плана до оптимизации
Связанные уроки:
Симптомы
- Scalar UDF не компилируется с ошибкой trait bound
- Попытка использовать `as_primitive_array::<Utf8>()` вызывает ошибку типов
- Rust не находит реализацию ArrowPrimitiveType для строковых типов
Причина
В Arrow типы Utf8/LargeUtf8/Binary — не примитивные. Они хранятся как offset-буфер + данные, а не как массив значений фиксированного размера. Для строковых массивов нужны специализированные downcast-методы.
Решение
- Используйте
as_string::<i32>(array)вместоas_primitive_array::<Utf8>(array) - Для бинарных данных:
as_binary::<i32>(array)илиas_large_binary(array) - Изучите набор downcast-функций в
arrow::array— каждый тип данных имеет свой метод
Связанные уроки:
Симптомы
- UDF регистрируется успешно, но падает при вызове в запросе
- Ошибка 'returned a different schema' при выполнении SQL с пользовательской функцией
- Тип возврата UDF не совпадает с объявленным в return_type()
Причина
Функция return_type() в ScalarUDF объявляет один тип, а фактический ArrayRef, возвращаемый из invoke(), содержит другой. DataFusion проверяет соответствие после выполнения функции.
Решение
- Сравните тип в
return_type()с типом ArrayRef вinvoke() - Если тип зависит от входных данных, используйте
return_type_from_exprs()вместо статического return_type - Используйте
arrow::compute::cast()в конце invoke() для приведения к ожидаемому типу
Связанные уроки:
Симптомы
- Кастомный TableProvider падает с паникой при первом запросе
- SELECT * FROM custom_table вызывает panic в рантайме
- Схема, возвращаемая scan(), не совпадает со schema()
Причина
DataFusion валидирует, что RecordBatch из ExecutionPlan (возвращённого scan()) соответствует схеме, объявленной в schema(). Если TableProvider создаёт схему динамически или с ошибкой, возникает несовпадение.
Решение
- Убедитесь, что
schema()иscan()используют один и тот же объект Arc<Schema> - Сохраните схему в поле структуры при создании провайдера:
self.schema.clone() - Если schema зависит от проекции — проверьте, что
scan(projection)корректно фильтрует поля
Связанные уроки:
Симптомы
- Пользовательское правило оптимизации вызывает панику
- Запросы с определёнными паттернами (JOIN, подзапрос) крашат оптимизатор
- Panic с трейсом, указывающим на кастомный OptimizerRule
Причина
Кастомный OptimizerRule неправильно обрабатывает структуру логического плана. Типичные причины: обращение к children по индексу без проверки длины, отсутствие обработки всех вариантов LogicalPlan enum, или мутация плана без пересчёта схемы.
Решение
- Оборачивайте тело правила в
catch_unwindна этапе отладки для получения backtrace - Проверяйте plan.inputs().len() перед обращением по индексу
- Используйте
plan.map_children()илиplan.rewrite()вместо ручного обхода дерева - Покройте правило тестами с EXPLAIN для различных паттернов запросов
Связанные уроки:
Симптомы
- Проект не компилируется после добавления зависимости datafusion
- cargo build выдаёт ошибку 'use of undeclared crate or module'
- Ошибка появляется при импорте datafusion::prelude::*
Причина
Версия datafusion в Cargo.toml не совпадает с версиями смежных крейтов (datafusion-common, datafusion-expr, arrow). DataFusion требует строгого совпадения версий всех своих подкрейтов.
Решение
- Используйте одну версию для всех datafusion-* крейтов в Cargo.toml
- Проверьте совместимость:
cargo tree -i datafusion— все версии должны совпадать - Для DataFusion 40+ используйте
arrow = { version = "53" }(версия Arrow привязана к конкретной версии DataFusion)
Связанные уроки:
Симптомы
- JOIN-запрос между большими таблицами выполняется медленно
- EXPLAIN показывает неоптимальный порядок JOIN (большая таблица слева)
- Планировщик не использует hash join, когда это было бы эффективнее
Причина
Без статистики (row count, column cardinality) оптимизатор не может выбрать порядок JOIN и тип JOIN. DataFusion использует эвристики, которые часто неоптимальны для больших данных.
Решение
- Зарегистрируйте таблицу с включённой статистикой:
CsvReadOptions::new().has_header(true)+ANALYZE TABLE - Для Parquet файлов статистика считывается автоматически из row group metadata
- Вручную укажите hint:
SELECT /*+ HASH_JOIN(a, b) */ ...(если поддерживается) - Задайте порядок вручную: маленькую таблицу ставьте в правую сторону JOIN
Связанные уроки:
Симптомы
- Hash Join падает с out-of-memory на больших таблицах
- Запрос с несколькими JOIN потребляет всю доступную память
- Декартово произведение из-за отсутствия или неправильного условия JOIN
Причина
Hash Join строит хеш-таблицу целиком в памяти для правой стороны JOIN. Если правая таблица огромна или условие JOIN слишком широкое (cross join), происходит взрыв памяти.
Решение
- Проверьте условие JOIN — убедитесь, что есть ON clause и он достаточно селективен
- Поместите меньшую таблицу справа:
big_table JOIN small_table - Включите spill-to-disk: настройте
FairSpillPoolс путём для временных файлов - Разбейте запрос: материализуйте промежуточные результаты через CREATE TABLE AS
Связанные уроки:
Симптомы
- Запрос использует только одно ядро CPU при наличии многих
- EXPLAIN ANALYZE показывает partitions=1 для большого файла
- Время выполнения не улучшается при увеличении target_partitions
Причина
Один большой файл (CSV или JSON) читается как одна партиция, так как DataFusion не может разбить текстовый файл на параллельные чанки. Parquet файлы с одним row group имеют ту же проблему.
Решение
- Разбейте данные на несколько файлов и зарегистрируйте директорию:
ctx.register_parquet("t", "data/", ...) - Конвертируйте CSV в Parquet с несколькими row groups
- Настройте
target_partitionsв SessionConfig:config.with_target_partitions(num_cpus::get()) - Добавьте repartition после чтения:
.repartition(Partitioning::RoundRobinBatch(8))
Связанные уроки:
Симптомы
- import datafusion вызывает ModuleNotFoundError
- pip install datafusion завершился успешно, но модуль не импортируется
- В virtualenv модуль не найден
Причина
Пакет datafusion-python требует совместимой версии Python (3.8+) и может не собираться на некоторых платформах (ARM Linux). Также возможен конфликт между системным Python и virtualenv.
Решение
- Установите в активный virtualenv:
python -m venv .venv && source .venv/bin/activate && pip install datafusion - Проверьте установку:
pip show datafusion— убедитесь, что Location соответствует sys.path - На macOS ARM:
pip install datafusion --no-binary :all:для сборки из исходников - Убедитесь, что используете Python 3.8+:
python --version
Связанные уроки:
Симптомы
- Ошибка при передаче PyArrow RecordBatch в DataFusion Python
- ArrowInvalid при вызове ctx.register_record_batches()
- После обновления PyArrow старый код перестал работать
Причина
datafusion-python привязан к конкретной версии PyArrow через Rust-биндинги (pyo3-arrow). Если установленная версия PyArrow не совпадает с ожидаемой, данные передаются с неправильным маппингом типов.
Решение
- Проверьте совместимость:
pip show datafusion pyarrow— версии должны быть совместимы - Установите рекомендуемую версию PyArrow:
pip install 'pyarrow>=14.0,<16.0'(для datafusion 37+) - Пересоздайте virtualenv с чистыми зависимостями:
pip install datafusionподтянет совместимый pyarrow
Связанные уроки:
Симптомы
- Передача pandas DataFrame в ctx.register_record_batches() вызывает TypeError
- Попытка использовать pandas объекты напрямую в DataFusion API
- Ошибка при смешивании pandas и Arrow API
Причина
DataFusion Python работает с PyArrow, а не с pandas напрямую. Функции регистрации ожидают RecordBatch или RecordBatchReader. Pandas DataFrame нужно конвертировать через PyArrow.
Решение
- Конвертируйте DataFrame в Arrow:
table = pa.Table.from_pandas(df)затемctx.register_record_batches("t", [table.to_batches()]) - Или используйте
ctx.from_pandas(df)если такой метод доступен в вашей версии - Для обратной конвертации:
result.to_pandas()на результате запроса
Связанные уроки:
Симптомы
- Ballista клиент не может подключиться к scheduler
- BallistaContext::remote() возвращает Connection refused
- Scheduler процесс запущен, но не принимает соединения
Причина
Scheduler не запущен, не успел инициализироваться, или слушает на другом адресе/порту. В Docker-окружении возможна проблема с сетью между контейнерами.
Решение
- Проверьте статус scheduler:
curl http://localhost:50050/api/state - Убедитесь, что scheduler слушает:
ss -tlnp | grep 50050 - В Docker: используйте имя сервиса вместо localhost:
BallistaContext::remote("scheduler", 50050, ...) - Дождитесь инициализации: scheduler может запускаться 5-10 секунд
Связанные уроки:
Симптомы
- Запрос с кастомным UDF выполняется локально, но падает в Ballista
- Ошибка сериализации при отправке плана на executor
- UDF зарегистрирована в клиенте, но не найдена на worker
Причина
В распределённом режиме физический план сериализуется и отправляется на executor-ы. Кастомные UDF должны быть зарегистрированы на каждом executor при старте, иначе десериализация плана не найдёт функцию.
Решение
- Зарегистрируйте UDF в FunctionRegistry на каждом executor при старте
- Используйте
--udf-libфлаг для загрузки shared library с UDF - Альтернатива: замените UDF на стандартные SQL-выражения для распределённых запросов
Связанные уроки:
Симптомы
- Распределённый запрос зависает на фазе shuffle
- Некоторые executor-ы не могут обменяться данными
- Таймауты при чтении промежуточных результатов
Причина
Executor-ы не могут связаться друг с другом напрямую (peer-to-peer). Частые причины: firewall блокирует порты между нодами, неправильная конфигурация внешнего адреса executor (advertise address), или сеть Docker overlay не настроена.
Решение
- Убедитесь, что порты executor-ов доступны между нодами (по умолчанию 50051)
- Задайте внешний адрес:
--external-host <node_ip>при запуске executor - В Docker Compose: используйте единую overlay-сеть для всех сервисов
- Увеличьте таймауты:
--shuffle-reader-timeout 60s
Связанные уроки:
Симптомы
- Docker build зависает и крашится при компиляции datafusion
- Процесс cargo build убивается с signal 9 (SIGKILL) внутри контейнера
- Docker Desktop показывает 100% использование памяти
Причина
Компиляция DataFusion и Arrow из исходников требует 4-8 GB RAM. Docker Desktop по умолчанию ограничивает контейнеры 2 GB, что недостаточно для параллельной компиляции Rust.
Решение
- Увеличьте лимит памяти Docker Desktop до 8 GB (Settings → Resources → Memory)
- Ограничьте параллельность:
ENV CARGO_BUILD_JOBS=2в Dockerfile - Используйте multi-stage build: собирайте в builder-образе, копируйте только бинарник в финальный
- Используйте готовые бинарные образы:
FROM datafusion/datafusion-cli:latest
Связанные уроки:
Симптомы
- Запись результатов запроса в файл внутри контейнера падает с Permission denied
- Volume mount не имеет прав на запись
- Контейнер запускается от root, но volume принадлежит другому пользователю
Причина
Docker volume mount наследует права хост-директории. Если процесс внутри контейнера работает от другого UID, чем владелец директории на хосте, запись невозможна.
Решение
- Задайте UID при запуске:
docker run -u $(id -u):$(id -g) ... - Создайте директорию заранее с правами:
mkdir -p ./output && chmod 777 ./output - В Dockerfile:
RUN mkdir /data/output && chown 1000:1000 /data/output - Используйте named volumes вместо bind mounts:
docker volume create datafusion_output
Связанные уроки:
Симптомы
- Чтение Parquet файла падает с ошибкой schema mismatch
- Файл читался раньше, но после добавления колонки — ошибка
- Несколько Parquet файлов в директории имеют разные схемы
Причина
DataFusion при чтении директории Parquet файлов ожидает единую схему. Если файлы были созданы с разными версиями схемы (schema evolution), или первый файл задаёт схему, которой нет в остальных, возникает ошибка.
Решение
- Используйте
schema_mergeпри регистрации:ParquetReadOptions::default().schema_merge(true) - Укажите явную схему при чтении:
ctx.register_parquet_with_schema("t", path, schema, opts) - Перепишите файлы с единой схемой через DataFusion:
CREATE TABLE ... AS SELECT * FROM old_table