ANSI Mode: строгая типизация, implicit cast и совместимость
spark.sql.ansi.enabled включает стандартное SQL-поведение: запрещает silent overflow, implicit type coercion и деление на ноль. Разбираем что сломается при миграции и как это настроить.
Философия: ANSI SQL vs Hive Legacy¶
Когда вы впервые запускаете Spark SQL и делаете SELECT 1 / 0, результат вас удивит: Spark вернёт null вместо ошибки. Попробуйте разделить строку '123' на число - Spark тихо преобразует её в число и вернёт результат. Прибавьте единицу к максимальному значению INT (2 147 483 647) - получите отрицательное число без единого предупреждения.
Это не случайность и не баг. Это осознанное архитектурное решение, унаследованное от Apache Hive, на котором изначально строился Spark SQL.
История: почему Hive выбрал мягкое поведение¶
Apache Hive проектировался в эпоху, когда большие данные только входили в обиход. Аналитики переносили SQL-запросы из реляционных СУБД в Hadoop, и данные были грязными: смешанные типы, нули вместо значений, строки там, где ожидались числа. Жёсткое SQL-поведение порождало бы тысячи ошибок при каждом запуске на нечищеных данных.
Hive выбрал путь наименьшего сопротивления: при любой неоднозначности возвращай null или делай неявное преобразование. Запрос не падает - аналитик получает хоть какой-то результат и может разбираться дальше. Это было прагматичным решением для своего времени.
Spark унаследовал это поведение через совместимость с HiveQL. Результат - legacy Hive-режим, в котором Spark работает по умолчанию и сегодня.
Проблема: silent data corruption¶
Проблема с мягким поведением - оно скрывает баги. Рассмотрим реальный сценарий:
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit
spark = SparkSession.builder.getOrCreate()
sales_df = spark.createDataFrame([
(1, 100, 0),
(2, 200, 5),
(3, 300, 10),
], ["order_id", "revenue", "clicks"])
result = sales_df.withColumn("cpc", col("revenue") / col("clicks"))
result.show()
Вывод:
+--------+-------+------+----+
|order_id|revenue|clicks| cpc|
+--------+-------+------+----+
| 1| 100| 0|null|
| 2| 200| 5|40.0|
| 3| 300| 10|30.0|
+--------+-------+------+----+
На первый взгляд всё нормально: у первого заказа null в поле cpc. Но представьте, что этот результат уходит в агрегацию:
result.agg({"cpc": "avg"}).show()
Вывод:
+--------+
|avg(cpc)|
+--------+
| 35.0|
+--------+
null игнорируется при агрегации, и средняя стоимость клика оказывается 35, а не 30. Один заказ с нулевыми кликами тихо исключился из расчёта, и никто не заметил. Это классический silent data corruption - данные выглядят корректными, но математически неверны.
А теперь более коварный пример с overflow:
campaign_df = spark.createDataFrame([
(1, 2000000000, 500000000),
], ["campaign_id", "impressions", "extra"])
result = campaign_df.withColumn(
"total",
col("impressions").cast("int") + col("extra").cast("int")
)
result.show()
Вывод:
+-----------+-----------+------+-----------+
|campaign_id|impressions| extra| total|
+-----------+-----------+------+-----------+
| 1| 2000000000|500000000|-1794967296|
+-----------+-----------+------+-----------+
Результат - отрицательное число! Сумма двух положительных чисел дала отрицательный результат из-за integer overflow, и Spark не выдал ни одного предупреждения. Этот результат может уйти в отчёт, в BI-дашборд, в бизнес-решение.
ANSI SQL стандарт и его цель¶
ANSI (American National Standards Institute) и ISO разрабатывают стандарты SQL с 1987 года. Ключевой принцип стандарта: SQL должен быть предсказуемым и детерминированным. Если операция математически некорректна (деление на ноль), запрос должен завершиться ошибкой, а не возвращать неопределённый результат.
Современные СУБД (PostgreSQL, MySQL в строгом режиме, SQL Server, Oracle) следуют ANSI-стандарту именно по этой причине: разработчик должен знать о проблеме и явно её обработать, а не получать тихо искажённые данные.
ANSI mode в Spark (spark.sql.ansi.enabled = true) - это переключатель, который приближает Spark SQL к поведению стандартных реляционных СУБД. Вместо тихого возврата null или переполнения - явные исключения. Вместо неявных приведений типов - требование явного CAST.
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.sql("SELECT 1 / 0").show()
java.lang.ArithmeticException: [DIVIDE_BY_ZERO]
Division by zero. Use try_divide to tolerate divisor being 0
and return NULL instead. If necessary set "spark.sql.ansi.enabled"
to "false" to bypass this error.
; 'Project [unresolvedstar()]
+- 'Project [(1 / 0) AS (1 / 0)#0]
+- OneRowRelation
Теперь ошибка явна, с чётким сообщением и подсказкой, как её обработать. Именно такое поведение - норма для промышленных SQL-систем.
Как включить ANSI mode¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("ANSI Example") \
.config("spark.sql.ansi.enabled", "true") \
.getOrCreate()
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.sql("SET spark.sql.ansi.enabled=true")
Все три варианта эквивалентны. При старте сессии - первый или третий. Для изменения в runtime - второй или третий.
Карта изменений: что меняется в ANSI mode¶
Диаграмма показывает шесть классов операций, которые Spark обрабатывает по-разному в зависимости от режима. В Legacy-режиме (левая ветка) все ошибки "поглощаются" - запрос продолжает работу, возвращая null или некорректные данные. В ANSI-режиме (правая ветка) каждая операция порождает конкретный тип исключения, который можно обработать или предотвратить.
Деление на ноль¶
Деление на ноль - математически неопределённая операция. В математике она не имеет результата. В SQL стандарт требует ошибки. Spark в Legacy-режиме нарушает этот принцип, возвращая null.
Поведение в Legacy и ANSI режимах¶
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("SELECT 10 / 0").show()
+---------+
|(10 / 0) |
+---------+
| null|
+---------+
Spark не только не падает - он даже не предупреждает. Если null попадёт в SUM() или AVG(), он будет проигнорирован, что может исказить агрегаты.
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
spark.sql("SELECT 10 / 0").show()
except Exception as e:
print(type(e).__name__, str(e)[:120])
AnalysisException [DIVIDE_BY_ZERO] Division by zero. Use try_divide to tolerate divisor being 0 and return NULL instead.
Деление на ноль с переменными¶
Деление на константный ноль Spark может поймать на этапе анализа. Но деление на столбец - только при выполнении:
from pyspark.sql.functions import col
data = [(1, 100, 5), (2, 200, 0), (3, 300, 10)]
df = spark.createDataFrame(data, ["id", "revenue", "sessions"])
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
df.withColumn("rps", col("revenue") / col("sessions")).show()
except Exception as e:
print(type(e).__name__, str(e)[:150])
SparkArithmeticException [DIVIDE_BY_ZERO] Division by zero. Use try_divide to tolerate divisor being 0 and return NULL instead.
Ошибка возникает при выполнении, когда Spark встречает строку с sessions = 0. Задача становится Failed, и весь DataFrame не вычисляется.
Функция try_divide¶
try_divide - безопасная обёртка деления, которая возвращает null при делении на ноль вместо исключения. Это "разрешённый" способ работы с потенциально нулевым делителем в ANSI-режиме:
from pyspark.sql.functions import try_divide
spark.conf.set("spark.sql.ansi.enabled", "true")
df.withColumn("rps", try_divide(col("revenue"), col("sessions"))).show()
+---+-------+--------+----+
| id|revenue|sessions| rps|
+---+-------+--------+----+
| 1| 100| 5|20.0|
| 2| 200| 0|null|
| 3| 300| 10|30.0|
+---+-------+--------+----+
Второй аргумент try_divide можно задать явно - это значение, которое возвращается при делении на ноль вместо null:
from pyspark.sql.functions import try_divide, lit
df.withColumn("rps", try_divide(col("revenue"), col("sessions"), lit(0.0))).show()
+---+-------+--------+----+
| id|revenue|sessions| rps|
+---+-------+--------+----+
| 1| 100| 5|20.0|
| 2| 200| 0| 0.0|
| 3| 300| 10|30.0|
+---+-------+--------+----+
CASE WHEN как альтернатива¶
До появления try_divide (Spark 3.3+) единственным способом защиты была явная проверка:
from pyspark.sql.functions import when, col
df.withColumn(
"rps",
when(col("sessions") > 0, col("revenue") / col("sessions"))
.otherwise(None)
).show()
Этот паттерн работает и в Legacy, и в ANSI режимах, и совместим со старыми версиями Spark. Однако try_divide предпочтительнее - он короче, выразительнее и явно документирует намерение.
Деление в SQL-запросах¶
-- Legacy: не упадёт
SELECT order_id, amount / quantity AS unit_price FROM orders;
-- ANSI safe: явная защита
SELECT
order_id,
CASE WHEN quantity > 0
THEN amount / quantity
ELSE NULL
END AS unit_price
FROM orders;
-- ANSI safe: через try_divide (Spark 3.3+)
SELECT order_id, try_divide(amount, quantity) AS unit_price FROM orders;
Арифметическое переполнение (Integer Overflow)¶
Integer overflow - ситуация, когда результат арифметической операции выходит за пределы диапазона типа данных. В двоичной арифметике это приводит к "перекручиванию": максимальное значение + 1 = минимальное значение.
Диапазоны целочисленных типов в Spark¶
| Тип | Байт | Минимум | Максимум |
|---|---|---|---|
ByteType |
1 | -128 | 127 |
ShortType |
2 | -32 768 | 32 767 |
IntegerType |
4 | -2 147 483 648 | 2 147 483 647 |
LongType |
8 | -9 223 372 036 854 775 808 | 9 223 372 036 854 775 807 |
Для INT (IntegerType) максимум - около 2.1 миллиарда. Для современных задач (количество просмотров, клики, байты данных) этого может не хватать.
Демонстрация overflow¶
from pyspark.sql.functions import col, lit
spark.conf.set("spark.sql.ansi.enabled", "false")
overflow_df = spark.createDataFrame([(2147483647,)], ["max_int"])
result = overflow_df.withColumn(
"overflow",
col("max_int") + lit(1)
)
result.show()
result.printSchema()
+-----------+-----------+
| max_int| overflow|
+-----------+-----------+
| 2147483647|-2147483648|
+-----------+-----------+
root
|-- max_int: integer (nullable = true)
|-- overflow: integer (nullable = true)
В Legacy-режиме результат -2147483648 - полностью некорректный. Запрос выполнился успешно, никаких предупреждений нет.
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
result = overflow_df.withColumn("overflow", col("max_int") + lit(1))
result.show()
except Exception as e:
print(type(e).__name__, ":", str(e)[:100])
SparkArithmeticException : [INTEGER_OVERFLOW] integer overflow in 'Add', values: 2147483647, 1.
В ANSI-режиме ошибка явная: переполнение IntegerType при операции сложения с конкретными значениями.
try_add, try_subtract, try_multiply¶
В Spark 3.3+ появились функции-аналоги для всех арифметических операций:
from pyspark.sql.functions import try_add, try_subtract, try_multiply
spark.conf.set("spark.sql.ansi.enabled", "true")
result = overflow_df.select(
try_add(col("max_int"), lit(1)).alias("safe_add"),
try_subtract(lit(-2147483648), lit(1)).alias("safe_sub"),
try_multiply(col("max_int"), lit(2)).alias("safe_mul"),
)
result.show()
+--------+--------+--------+
|safe_add|safe_sub|safe_mul|
+--------+--------+--------+
| null| null| null|
+--------+--------+--------+
При переполнении функции возвращают null вместо исключения или некорректного значения. Это позволяет продолжить вычисления и обработать исключительные случаи явно.
Правильный подход: выбор типа¶
Лучший способ избежать overflow - выбирать тип с запасом:
from pyspark.sql.functions import col
df.withColumn(
"total_clicks",
col("daily_clicks").cast("long") * lit(365)
)
Приведение к LongType перед операцией устраняет риск переполнения для большинства практических сценариев (LongType вмещает до 9.2 × 10^18). Для финансовых расчётов с дробными значениями используйте DecimalType:
from pyspark.sql.types import DecimalType
df.withColumn(
"revenue",
col("amount").cast(DecimalType(18, 4))
)
DecimalType(18, 4) означает: 18 значащих цифр всего, 4 после запятой. При переполнении DecimalType в ANSI-режиме также выбрасывает исключение.
Неявное приведение типов (Implicit Type Coercion)¶
Type coercion - автоматическое преобразование типов данных, когда операнды имеют несовместимые типы. Spark Legacy-режим делает это тихо и по своим правилам, что может приводить к неожиданным результатам.
Иерархия автоматических преобразований в Legacy¶
В Legacy-режиме Spark следует иерархии: ByteType < ShortType < IntegerType < LongType < FloatType < DoubleType. При смешивании типов меньший преобразуется в больший. Строки имеют особый статус - Spark пытается их распарсить.
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("SELECT '123' + 456").show()
spark.sql("SELECT true + 1").show()
spark.sql("SELECT 1.5 + 2").show()
spark.sql("SELECT '2024-01-01' + INTERVAL 1 DAY").show()
+-----------+
|(123 + 456)|
+-----------+
| 579|
+-----------+
+----------+
|(true + 1)|
+----------+
| 2|
+----------+
+---------+
|(1.5 + 2)|
+---------+
| 3.5|
+---------+
+---------------------------------------------+
|(CAST(2024-01-01 AS DATE) + INTERVAL '1' DAY)|
+---------------------------------------------+
| 2024-01-02|
+---------------------------------------------+
Spark успешно преобразовал строку в число, boolean в целое, целое в дробное, и строку в дату. Всё это происходит автоматически, без вашего ведома.
Почему implicit coercion опасна¶
Проблема не в том, что преобразования происходят, а в том, что они непредсказуемы:
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("""
SELECT
'3.14' + 1 AS str_plus_int,
'abc' + 1 AS bad_str_plus_int
""").show()
+------------+----------------+
|str_plus_int|bad_str_plus_int|
+------------+----------------+
| 4.14| null|
+------------+----------------+
'3.14' + 1 даёт 4.14 (Spark распарсил строку как число), а 'abc' + 1 даёт null (распарсить не удалось, тихий null). Никакой ошибки нет - разработчик должен знать, что 'abc' не будет преобразован, и null тихо распространится дальше.
ANSI режим требует явного cast¶
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
spark.sql("SELECT '123' + 456").show()
except Exception as e:
print(str(e)[:150])
[DATATYPE_MISMATCH.BINARY_OP_WRONG_TYPE] Cannot resolve '(CAST('123' AS DOUBLE) + 456)'
due to data type mismatch: argument 1 requires (numeric or interval day to second or
interval year to month or interval) type, however, '123' is of string type.
Теперь нужно явно указать, что именно вы хотите:
spark.sql("SELECT CAST('123' AS INT) + 456").show()
+------------------+
|(CAST(123 AS INT) + 456)|
+------------------+
| 579|
+------------------+
Таблица допустимых преобразований в ANSI mode¶
В ANSI mode не все преобразования запрещены - только небезопасные:
| Исходный тип | Целевой тип | Legacy | ANSI |
|---|---|---|---|
IntegerType |
LongType |
Автоматически | Автоматически (widening) |
FloatType |
DoubleType |
Автоматически | Автоматически (widening) |
IntegerType |
DoubleType |
Автоматически | Автоматически (widening) |
LongType |
IntegerType |
Автоматически (сужение!) | Ошибка (потеря точности) |
StringType |
IntegerType |
Автоматически | Ошибка, нужен явный CAST |
BooleanType |
IntegerType |
Автоматически | Ошибка, нужен явный CAST |
DateType |
TimestampType |
Автоматически | Автоматически |
StringType |
DateType |
Автоматически | Ошибка, нужен явный CAST |
Расширяющие преобразования (widening - от меньшего типа к большему) разрешены в обоих режимах, потому что не теряют точность. Сужающие - запрещены в ANSI mode.
Boolean arithmetic¶
В SQL стандарт не разрешает использовать boolean в арифметике напрямую. В Legacy Spark это работает:
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("SELECT true + 1, false + 1, true * 5").show()
+----------+-----------+--------+
|(true + 1)|(false + 1)|(true*5)|
+----------+-----------+--------+
| 2| 1| 5|
+----------+-----------+--------+
true = 1, false = 0. В ANSI mode это запрещено:
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
spark.sql("SELECT true + 1").show()
except Exception as e:
print(str(e)[:100])
[DATATYPE_MISMATCH.BINARY_OP_WRONG_TYPE] Cannot resolve '(true + 1)' due to data type mismatch
Правильный путь - явный cast:
spark.sql("SELECT CAST(true AS INT) + 1").show()
spark.sql("SELECT IF(flag, 1, 0) + base_count FROM events").show()
try_cast: безопасное преобразование типов¶
CAST в ANSI mode выбрасывает исключение при невозможности преобразовать значение. TRY_CAST - его безопасный аналог, возвращающий null при неудаче.
Разница CAST и TRY_CAST¶
spark.conf.set("spark.sql.ansi.enabled", "true")
data = [("123",), ("abc",), ("456",), (None,), ("99.9",)]
df = spark.createDataFrame(data, ["raw_value"])
try:
df.select("raw_value", col("raw_value").cast("int").alias("as_int")).show()
except Exception as e:
print("CAST failed:", str(e)[:80])
df.select(
"raw_value",
col("raw_value").cast("int").alias("cast_int"),
).show()
В Legacy режиме CAST('abc' AS INT) вернёт null. В ANSI режиме это исключение. Используйте try_cast:
from pyspark.sql.functions import try_cast
df.select(
"raw_value",
try_cast(col("raw_value"), "int").alias("safe_int")
).show()
+---------+--------+
|raw_value|safe_int|
+---------+--------+
| 123| 123|
| abc| null|
| 456| 456|
| null| null|
| 99.9| null|
+---------+--------+
99.9 не преобразовалось в int (нет дробной части), поэтому null. Это корректно: строка не является валидным целым числом.
try_cast для дат¶
Особенно полезен try_cast при парсинге дат из сырых данных:
from pyspark.sql.functions import try_cast, col
dates_df = spark.createDataFrame([
("2024-01-15",),
("2024/02/30",),
("invalid-date",),
("2023-12-31",),
("",),
], ["raw_date"])
spark.conf.set("spark.sql.ansi.enabled", "true")
dates_df.select(
"raw_date",
try_cast(col("raw_date"), "date").alias("parsed_date")
).show()
+------------+-----------+
| raw_date|parsed_date|
+------------+-----------+
| 2024-01-15| 2024-01-15|
| 2024/02/30| null|
|invalid-date| null|
| 2023-12-31| 2023-12-31|
| | null|
+------------+-----------+
2024/02/30 не существует (30 февраля 2024 года), поэтому null. Это правильно и важно: в Legacy Spark может распарсить такую дату некорректно.
try_cast в SQL-запросах¶
-- SQL-синтаксис TRY_CAST
SELECT
user_id,
TRY_CAST(age_raw AS INT) AS age,
TRY_CAST(signup_date AS DATE) AS signup,
TRY_CAST(revenue_str AS DECIMAL(12,2)) AS revenue
FROM raw_events
WHERE TRY_CAST(age_raw AS INT) BETWEEN 18 AND 100;
TRY_CAST в WHERE позволяет фильтровать строки с непарсируемыми значениями без ошибок - они просто не пройдут условие BETWEEN.
Строки и VARCHAR: усечение¶
Тип VARCHAR(n) ограничивает длину строки. При попытке записать строку длиннее n символов поведение различается:
Поведение при усечении¶
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("""
CREATE OR REPLACE TEMP VIEW products AS
SELECT 'Very Long Product Name That Exceeds Limit' AS name
""")
spark.sql("""
SELECT CAST(name AS VARCHAR(10)) AS short_name FROM products
""").show()
+----------+
|short_name|
+----------+
|Very Long |
+----------+
В Legacy режиме строка тихо усечена до 10 символов. В ANSI режиме:
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
spark.sql("""
SELECT CAST('Very Long Product Name' AS VARCHAR(5)) AS s
""").show()
except Exception as e:
print(str(e)[:120])
[CAST_INVALID_INPUT] The value 'Very Long Product Name' of the type "STRING" cannot be cast to "VARCHAR(5)"
because it is malformed. Correct the value as per the syntax, or change its target type.
storeAssignmentPolicy¶
Конфигурация spark.sql.storeAssignmentPolicy управляет поведением при INSERT INTO с несовместимыми типами:
| Значение | Поведение |
|---|---|
ANSI |
Ошибка при несовместимости (рекомендуется) |
LEGACY |
Тихое приведение типов |
STRICT |
Строгая проверка, без implicit cast даже для widening |
spark.conf.set("spark.sql.storeAssignmentPolicy", "ANSI")
При ANSI политике попытка вставить строку 'toolong' в колонку VARCHAR(3) завершится ошибкой. Это защищает целостность данных в таблицах.
Array и Map: выход за границы¶
Операции с коллекциями в Spark также меняют поведение при ANSI mode.
Обращение к элементам массива¶
from pyspark.sql.functions import col
arr_df = spark.createDataFrame([
([1, 2, 3],),
([10, 20],),
], ["numbers"])
spark.conf.set("spark.sql.ansi.enabled", "false")
arr_df.select(col("numbers")[5]).show()
+----------+
|numbers[5]|
+----------+
| null|
| null|
+----------+
Индекс 5 выходит за пределы обоих массивов, но ошибки нет - null. В ANSI режиме:
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
arr_df.select(col("numbers")[5]).show()
except Exception as e:
print(str(e)[:100])
SparkArrayIndexOutOfBoundsException: [ARRAY_INDEX_OUT_OF_BOUNDS] The index 5 is out of bounds.
Array has 3 elements (indices from 0 to 2). Use the try_element_at function to avoid this error.
Сообщение содержит подсказку: try_element_at.
try_element_at¶
from pyspark.sql.functions import try_element_at
spark.conf.set("spark.sql.ansi.enabled", "true")
arr_df.select(
"numbers",
try_element_at(col("numbers"), lit(1)).alias("first"),
try_element_at(col("numbers"), lit(5)).alias("fifth"),
).show()
+---------+-----+-----+
| numbers|first|fifth|
+---------+-----+-----+
|[1, 2, 3]| 1| null|
| [10, 20]| 10| null|
+---------+-----+-----+
Внимание: в Spark element_at использует 1-based индексацию (первый элемент - 1, а не 0). Отрицательные индексы считают с конца: -1 - последний элемент.
Map: ключи, которых нет¶
from pyspark.sql.functions import try_element_at, lit
map_df = spark.createDataFrame([
({"a": 1, "b": 2},),
], ["data"])
spark.conf.set("spark.sql.ansi.enabled", "true")
map_df.select(
try_element_at(col("data"), lit("a")).alias("key_a"),
try_element_at(col("data"), lit("z")).alias("key_z"),
).show()
+-----+-----+
|key_a|key_z|
+-----+-----+
| 1| null|
+-----+-----+
try_element_at для Map возвращает null при отсутствии ключа вместо исключения.
Строгий парсинг дат и времени¶
ANSI mode также влияет на парсинг временных меток.
Форматы дат: строгие vs мягкие¶
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("SELECT CAST('2024-13-01' AS DATE)").show()
spark.sql("SELECT CAST('2024-01-32' AS DATE)").show()
+----+
|null|
+----+
+----+
|null|
+----+
13-й месяц и 32-й день - несуществующие значения, но в Legacy они возвращают null.
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
spark.sql("SELECT CAST('2024-13-01' AS DATE)").show()
except Exception as e:
print(str(e)[:120])
[CAST_INVALID_INPUT] The value '2024-13-01' of the type "STRING" cannot be cast to "DATE"
because it is malformed.
Конфигурация временных зон и форматов¶
spark.conf.set("spark.sql.session.timeZone", "UTC")
spark.conf.set("spark.sql.datetime.java8API.enabled", "true")
java8API.enabled переключает Spark на java.time API вместо устаревшего java.util.Date. Это влияет на точность временных зон и обработку исторических дат.
to_date и to_timestamp с явным форматом¶
В обоих режимах рекомендуется задавать формат явно:
from pyspark.sql.functions import to_date, to_timestamp, col
events_df = spark.createDataFrame([
("15-01-2024 10:30:00",),
("32-01-2024 00:00:00",),
("invalid",),
], ["raw_ts"])
events_df.select(
"raw_ts",
to_date(col("raw_ts"), "dd-MM-yyyy").alias("event_date"),
to_timestamp(col("raw_ts"), "dd-MM-yyyy HH:mm:ss").alias("event_ts"),
).show()
+-------------------+----------+-------------------+
| raw_ts|event_date| event_ts|
+-------------------+----------+-------------------+
|15-01-2024 10:30:00|2024-01-15|2024-01-15 10:30:00|
|32-01-2024 00:00:00| null| null|
| invalid| null| null|
+-------------------+----------+-------------------+
to_date и to_timestamp возвращают null при неудаче независимо от режима ANSI. Они не выбрасывают исключения при неверном формате. Это их стандартное поведение - по сути, они всегда ведут себя как try_cast.
Зарезервированные ключевые слова¶
ANSI mode расширяет список зарезервированных слов SQL. Имена колонок или таблиц, совпадающие с зарезервированными словами, должны экранироваться обратными кавычками.
Слова, добавляемые в ANSI mode¶
spark.conf.set("spark.sql.ansi.enabled", "false")
spark.sql("""
SELECT 1 AS value, 2 AS table, 3 AS end
""").show()
+-----+-----+---+
|value|table|end|
+-----+-----+---+
| 1| 2| 3|
+-----+-----+---+
В Legacy режиме table и end как имена колонок - без проблем. В ANSI:
spark.conf.set("spark.sql.ansi.enabled", "true")
try:
spark.sql("SELECT 1 AS end").show()
except Exception as e:
print(str(e)[:150])
[PARSE_SYNTAX_ERROR] Syntax error at or near 'end'(line 1, col 12): SELECT 1 AS end
end - зарезервированное слово ANSI SQL. Экранирование обратными кавычками решает проблему:
spark.sql("SELECT 1 AS `end`").show()
+---+
|end|
+---+
| 1|
+---+
Список часто встречающихся зарезервированных слов¶
Полный список зарезервированных слов в ANSI SQL включает: ALL, ANY, ARRAY, AS, AT, BETWEEN, BOTH, BY, CASE, CAST, CHECK, COLUMN, CONSTRAINT, CREATE, CROSS, CUBE, CURRENT, CURRENT_DATE, CURRENT_TIME, CURRENT_TIMESTAMP, CURRENT_USER, DELETE, DISTINCT, DROP, ELSE, END, EXCEPT, EXISTS, FETCH, FILTER, FOLLOWING, FOR, FOREIGN, FROM, FULL, GROUPING, HAVING, IN, INNER, INSERT, INTERSECT, INTO, IS, JOIN, LEADING, LEFT, LIKE, LIMIT, LOCAL, NATURAL, NO, NOT, NULL, OF, ON, ORDER, OUTER, OVER, PARTITION, PERCENT, PRECEDING, RANGE, REFERENCES, RIGHT, ROLLUP, ROW, ROWS, SELECT, SESSION_USER, SET, SOME, START, TABLE, TIME, TO, TRAILING, UNION, UNIQUE, UPDATE, USER, USING, VALUE, VALUES, WHEN, WHERE, WINDOW, WITH.
Практический совет: если колонка называется value, table, end, from - обязательно экранируйте в ANSI mode или переименуйте при выборке.
-- Безопасный паттерн
SELECT
`from` AS source_date,
`table` AS table_name,
`end` AS end_date,
`value` AS metric_value
FROM legacy_table;
Совместимость с другими СУБД¶
Одна из главных причин включать ANSI mode - совместимость с другими системами, которые SQL-стандарт соблюдают строго.
Диаграмма совместимости¶
Диаграмма показывает экосистему: Spark с включённым ANSI mode становится совместимым партнёром для всех перечисленных систем. Каждая из них реализует строгое SQL-поведение, и запрос, написанный для PostgreSQL, должен работать в Spark с ANSI mode без модификаций.
PostgreSQL¶
PostgreSQL - самая близкая к ANSI SQL реляционная СУБД. Если вы переносите аналитику из PostgreSQL в Spark:
-- PostgreSQL: ошибка при делении на ноль
SELECT revenue / sessions FROM events WHERE campaign_id = 1;
-- ERROR: division by zero
-- Spark Legacy: null
-- Spark ANSI: то же поведение, что PostgreSQL - ошибка
-- Переносимый код для обеих систем:
SELECT CASE WHEN sessions > 0
THEN revenue / sessions
ELSE NULL
END AS cpc
FROM events WHERE campaign_id = 1;
Trino / Presto¶
Trino часто работает как федеративный движок поверх тех же данных, что и Spark. Запросы должны давать одинаковые результаты:
-- Trino: CAST('abc' AS INTEGER) → NULL (!) или ошибка зависит от настройки
-- Spark ANSI: исключение
-- Переносимый код:
SELECT TRY_CAST(raw_value AS INTEGER) AS parsed_int FROM raw_table;
-- TRY_CAST поддерживается в Trino, Spark 3.2+, Snowflake, BigQuery (TRY_CAST аналог - SAFE_CAST)
Snowflake¶
Snowflake использует термин TRY_TO_<TYPE> вместо TRY_CAST:
| Spark | Snowflake | BigQuery |
|---|---|---|
TRY_CAST(x AS INTEGER) |
TRY_TO_NUMBER(x) |
SAFE_CAST(x AS INT64) |
TRY_CAST(x AS DATE) |
TRY_TO_DATE(x) |
SAFE.PARSE_DATE(fmt, x) |
TRY_DIVIDE(a, b) |
IFF(b != 0, a/b, NULL) |
SAFE_DIVIDE(a, b) |
При разработке кросс-платформенных пайплайнов нужно знать эти соответствия. Задача ETL-разработчика - писать SQL, который работает корректно на всех платформах, или иметь чёткий слой адаптации.
BigQuery¶
BigQuery использует SAFE. префикс или SAFE_CAST/SAFE_DIVIDE для безопасных операций:
-- BigQuery
SELECT SAFE_DIVIDE(revenue, sessions) AS cpc FROM events;
SELECT SAFE_CAST(raw_date AS DATE) AS event_date FROM raw_events;
-- Spark ANSI эквивалент
SELECT TRY_DIVIDE(revenue, sessions) AS cpc FROM events;
SELECT TRY_CAST(raw_date AS DATE) AS event_date FROM raw_events;
Если вы пишете код для Spark и BigQuery одновременно - держите в уме эти синтаксические различия. Семантика одинакова, синтаксис - нет.
ANSI mode и Delta Lake / Apache Iceberg¶
Форматы таблиц с поддержкой ACID-транзакций также взаимодействуют с ANSI mode.
Delta Lake¶
Delta Lake наследует конфигурацию Spark - нет отдельного переключателя. Если spark.sql.ansi.enabled = true, то все операции с Delta таблицами (INSERT, UPDATE, MERGE) работают в строгом режиме:
from delta.tables import DeltaTable
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.sql("""
CREATE TABLE IF NOT EXISTS events_delta (
event_id BIGINT,
user_id INT,
amount DECIMAL(12, 4),
event_ts TIMESTAMP
) USING DELTA LOCATION '/data/events'
""")
spark.sql("""
INSERT INTO events_delta VALUES (1, 999999999999, 100.50, '2024-01-15 10:00:00')
""")
В ANSI mode вставка 999999999999 в колонку INT (максимум ~2.1 миллиарда) вызовет ошибку:
[INSERT_COLUMN_ARITY_MISMATCH.NOT_ENOUGH_DATA_COLUMNS] Cannot write to 'events_delta',
as the number of user-specified column names does not match the data schema.
Или, если тип не совпадает:
[STORE_ASSIGNMENT_VIOLATION] Assignment implicit type coercion is not allowed.
'user_id' requires int but 'expr' has type bigint.
Schema Evolution в ANSI mode¶
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.sql("""
ALTER TABLE events_delta
ADD COLUMNS (new_metric DOUBLE COMMENT 'New feature metric')
""")
spark.sql("""
ALTER TABLE events_delta
CHANGE COLUMN user_id user_id BIGINT
""")
Расширяющее изменение (INT → BIGINT) разрешено. Сужающее (BIGINT → INT) в ANSI mode может вызвать ошибку при чтении существующих данных.
Apache Iceberg¶
Iceberg имеет собственную систему типов, но при работе через Spark-коннектор наследует поведение ANSI:
spark.conf.set("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog")
spark.conf.set("spark.sql.catalog.iceberg.type", "hadoop")
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.sql("""
CREATE TABLE iceberg.db.events (
id BIGINT,
amount DECIMAL(18, 4)
) USING iceberg
""")
spark.sql("""
INSERT INTO iceberg.db.events VALUES (1, CAST('invalid' AS DECIMAL(18, 4)))
""")
В ANSI mode CAST('invalid' AS DECIMAL(18, 4)) - ошибка на этапе анализа. В Legacy - null в таблице.
Workflow миграции на ANSI mode¶
Переход существующего проекта на ANSI mode - это не кнопка "включить и забыть". Нужен системный подход, чтобы не сломать production.
Этап 1: Аудит - включить в dev/staging и собрать ошибки¶
spark.conf.set("spark.sql.ansi.enabled", "true")
import logging
logging.basicConfig(level=logging.WARNING)
logger = logging.getLogger("ansi_migration")
queries_to_test = [
("division_report", "SELECT revenue / sessions FROM events"),
("type_cast_etl", "SELECT '123' + visits FROM web_logs"),
("integer_math", "SELECT clicks * multiplier FROM campaigns"),
]
errors = []
for query_name, query in queries_to_test:
try:
spark.sql(query).limit(1).collect()
print(f"✅ {query_name}: OK")
except Exception as e:
errors.append((query_name, str(e)[:200]))
print(f"❌ {query_name}: {str(e)[:100]}")
print(f"\nTotal errors: {len(errors)}")
for name, err in errors:
print(f" {name}: {err}")
Запустите этот скрипт со всеми основными запросами вашего проекта. Собирайте не только исключения - проверяйте также результаты агрегаций (нет ли подозрительных null, нет ли нулей там, где раньше были значения).
Этап 2: Классификация ошибок¶
После сбора ошибок классифицируйте их:
| Тип ошибки | Причина | Решение |
|---|---|---|
DIVIDE_BY_ZERO |
Деление на нулевой столбец | try_divide, CASE WHEN |
INTEGER_OVERFLOW |
Тип слишком мал | Заменить на LONG, DECIMAL |
DATATYPE_MISMATCH |
Implicit coercion | Добавить явный CAST |
CAST_INVALID_INPUT |
Плохие данные при CAST | Заменить на TRY_CAST |
ARRAY_INDEX_OUT_OF_BOUNDS |
Доступ по индексу без проверки | try_element_at |
PARSE_SYNTAX_ERROR |
Зарезервированные слова | Экранирование ` |
Этап 3: Рефакторинг - типичные паттерны замены¶
Паттерн 1: Деление
df.withColumn("cpc", col("cost") / col("clicks"))
df.withColumn("cpc", try_divide(col("cost"), col("clicks")))
Паттерн 2: Implicit string cast
df.filter(col("date_col") > "2024-01-01")
df.filter(col("date_col") > lit("2024-01-01").cast("date"))
Паттерн 3: Integer operations
df.withColumn("total", col("count") * lit(365))
df.withColumn("total", col("count").cast("long") * lit(365L))
Паттерн 4: Boolean arithmetic
df.withColumn("flag_int", col("is_active") + lit(0))
df.withColumn("flag_int", col("is_active").cast("int"))
Паттерн 5: String to number
df.withColumn("amount", col("raw_amount").cast("double"))
df.withColumn("amount", try_cast(col("raw_amount"), "double"))
Этап 4: Постепенное включение - per-query approach¶
Если весь проект не готов к переходу сразу, включайте ANSI mode точечно:
def run_with_ansi(spark, query):
"""Запустить запрос в ANSI mode, восстановить исходный режим после."""
original = spark.conf.get("spark.sql.ansi.enabled", "false")
try:
spark.conf.set("spark.sql.ansi.enabled", "true")
return spark.sql(query)
finally:
spark.conf.set("spark.sql.ansi.enabled", original)
run_with_ansi(spark, "SELECT try_divide(revenue, sessions) FROM events").show()
Новые запросы и новые таблицы - в ANSI mode. Старый legacy-код - без изменений, до тех пор, пока не отрефакторен.
Этап 5: Автоматизированное тестирование¶
import pytest
from pyspark.sql import SparkSession
@pytest.fixture(scope="session")
def ansi_spark():
return SparkSession.builder \
.config("spark.sql.ansi.enabled", "true") \
.getOrCreate()
def test_division_safety(ansi_spark):
"""Проверить, что деление не падает при нулевых делителях."""
df = ansi_spark.createDataFrame(
[(1, 100, 0), (2, 200, 5)],
["id", "revenue", "clicks"]
)
result = df.withColumn(
"cpc",
try_divide(col("revenue"), col("clicks"))
)
rows = result.collect()
assert rows[0]["cpc"] is None
assert rows[1]["cpc"] == 40.0
def test_no_overflow(ansi_spark):
"""Проверить, что LONG используется там, где INT может переполниться."""
df = ansi_spark.createDataFrame([(2_000_000_000,)], ["large_int"])
result = df.withColumn(
"doubled",
col("large_int").cast("long") * 2
)
assert result.collect()[0]["doubled"] == 4_000_000_000
Конфигурация ANSI mode: полный справочник¶
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.conf.set("spark.sql.storeAssignmentPolicy", "ANSI")
spark.conf.set("spark.sql.storeAssignmentPolicy", "STRICT")
spark.conf.set("spark.sql.storeAssignmentPolicy", "LEGACY")
spark.conf.set("spark.sql.datetime.java8API.enabled", "true")
spark.conf.set("spark.sql.session.timeZone", "UTC")
Полная таблица конфигураций¶
| Конфигурация | По умолчанию | Описание |
|---|---|---|
spark.sql.ansi.enabled |
false |
Основной переключатель ANSI mode |
spark.sql.storeAssignmentPolicy |
ANSI (3.x) |
Политика приведения типов при INSERT |
spark.sql.datetime.java8API.enabled |
false |
Использовать java.time вместо java.util.Date |
spark.sql.session.timeZone |
JVM default | Временная зона для datetime операций |
spark.sql.legacy.castComplexTypesToString.enabled |
false |
Legacy поведение CAST коллекций в строку |
spark.sql.legacy.timeParserPolicy |
CORRECTED |
LEGACY / CORRECTED / EXCEPTION |
timeParserPolicy¶
Особенно важна timeParserPolicy:
LEGACY- старый парсер дат, иногда принимает невалидные форматыCORRECTED- исправленный парсер (умолчание в Spark 3.x)EXCEPTION- строгий режим: невалидная дата → исключение (максимальная строгость)
spark.conf.set("spark.sql.legacy.timeParserPolicy", "EXCEPTION")
В режиме EXCEPTION даже to_date("2024-13-01", "yyyy-MM-dd") вызовет исключение вместо null. Используйте try_to_date (если доступна) или обрабатывайте через try-except.
Практика: ETL pipeline с ANSI mode¶
Разберём реальный сценарий: ETL pipeline, читающий сырые события из JSON, валидирующий и загружающий в Delta таблицу.
Задача¶
Есть Bronze таблица с сырыми событиями:
{"event_id": "1", "user_id": "100", "amount": "29.99", "ts": "2024-01-15 10:30:00", "clicks": "5"}
{"event_id": "2", "user_id": "abc", "amount": "-1", "ts": "2024-13-01", "clicks": "0"}
{"event_id": "3", "user_id": "200", "amount": "invalid", "ts": "2024-01-20", "clicks": "10"}
Все поля - строки (JSON). Нужно конвертировать типы, валидировать и загрузить в Silver.
Решение с ANSI mode¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import (
col, try_cast, try_divide, current_timestamp,
when, lit, count
)
from pyspark.sql.types import LongType, DoubleType, DateType, IntegerType
spark = SparkSession.builder \
.appName("ANSI ETL Pipeline") \
.config("spark.sql.ansi.enabled", "true") \
.config("spark.sql.storeAssignmentPolicy", "ANSI") \
.getOrCreate()
bronze_df = spark.createDataFrame([
("1", "100", "29.99", "2024-01-15 10:30:00", "5"),
("2", "abc", "-1", "2024-13-01", "0"),
("3", "200", "invalid", "2024-01-20", "10"),
("4", "300", "150.00", "2024-01-21", "3"),
], ["event_id", "user_id", "amount", "ts", "clicks"])
silver_df = bronze_df.select(
try_cast(col("event_id"), LongType()).alias("event_id"),
try_cast(col("user_id"), IntegerType()).alias("user_id"),
try_cast(col("amount"), DoubleType()).alias("amount"),
try_cast(col("ts"), DateType()).alias("event_date"),
try_cast(col("clicks"), IntegerType()).alias("clicks"),
)
silver_with_cpc = silver_df.withColumn(
"avg_revenue_per_click",
try_divide(col("amount"), col("clicks").cast("double"))
)
good_rows = silver_with_cpc.filter(
col("event_id").isNotNull() &
col("user_id").isNotNull() &
col("amount").isNotNull() &
col("event_date").isNotNull()
)
bad_rows = silver_with_cpc.filter(
col("event_id").isNull() |
col("user_id").isNull() |
col("amount").isNull() |
col("event_date").isNull()
)
print("=== Good rows ===")
good_rows.show()
print("=== Bad rows (quarantine) ===")
bad_rows.show()
total = silver_with_cpc.count()
good = good_rows.count()
bad = bad_rows.count()
print(f"Quality report: {good}/{total} valid, {bad}/{total} quarantined ({bad/total*100:.1f}%)")
Вывод:
=== Good rows ===
+--------+-------+------+----------+------+---------------------+
|event_id|user_id|amount|event_date|clicks|avg_revenue_per_click|
+--------+-------+------+----------+------+---------------------+
| 1| 100| 29.99|2024-01-15| 5| 5.998|
| 4| 300|150.00|2024-01-21| 3| 50.000|
+--------+-------+------+----------+------+---------------------+
=== Bad rows (quarantine) ===
+--------+-------+------+----------+------+---------------------+
|event_id|user_id|amount|event_date|clicks|avg_revenue_per_click|
+--------+-------+------+----------+------+---------------------+
| 2| null| -1.0| null| 0| null|
| 3| 200| null|2024-01-20| 10| null|
+--------+-------+------+----------+------+---------------------+
Quality report: 2/4 valid, 2/4 quarantined (50.0%)
Строка 2: user_id = 'abc' не распарсился → null; ts = '2024-13-01' - несуществующая дата → null. Строка 3: amount = 'invalid' не распарсился → null. Все ошибки - через try_cast, без исключений. Строка 2 с clicks = 0 - try_divide вернул null вместо ошибки.
good_rows.write \
.format("delta") \
.mode("append") \
.save("/silver/events")
bad_rows.withColumn("quarantine_ts", current_timestamp()) \
.write \
.format("delta") \
.mode("append") \
.save("/silver/events_quarantine")
Антипаттерны¶
Антипаттерн 1: Включить ANSI mode без тестирования¶
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.sql("INSERT INTO production_table SELECT * FROM etl_result")
Без предварительного аудита production ETL упадёт в неожиданных местах. Всегда сначала тестируйте в dev/staging со всеми данными.
Антипаттерн 2: Заменить CAST на CAST вместо TRY_CAST¶
raw_df.withColumn("amount", col("raw").cast("double"))
В ANSI mode это упадёт при первой нечисловой строке. Если данные сырые - всегда try_cast:
raw_df.withColumn("amount", try_cast(col("raw"), "double"))
Антипаттерн 3: Не обрабатывать null после try_cast¶
df.withColumn("amount", try_cast(col("raw"), "double")) \
.agg({"amount": "avg"})
avg молча проигнорирует null. Добавьте мониторинг качества:
df_with_types = df.withColumn("amount", try_cast(col("raw"), "double"))
null_rate = df_with_types.filter(col("amount").isNull()).count() / df.count()
if null_rate > 0.05:
raise ValueError(f"Too many null amounts: {null_rate:.1%}")
df_with_types.agg({"amount": "avg"}).show()
Антипаттерн 4: Включить ANSI для всего кластера в spark-defaults.conf¶
# spark-defaults.conf
spark.sql.ansi.enabled=true # ОПАСНО для shared cluster
На shared-кластере это ломает все существующие приложения, которые полагаются на Legacy-поведение. Используйте per-application или per-session конфигурацию.
Антипаттерн 5: Boolean arithmetic без явного cast¶
df.withColumn("active_count", col("is_active") + col("total_count"))
В ANSI mode - ошибка. Явный cast:
df.withColumn("active_count", col("is_active").cast("int") + col("total_count"))
Антипаттерн 6: Игнорировать зарезервированные слова¶
spark.conf.set("spark.sql.ansi.enabled", "true")
df.createOrReplaceTempView("events")
spark.sql("SELECT end, value, table FROM events")
end, value, table - зарезервированные слова. Используйте `:
spark.sql("SELECT `end`, `value`, `table` FROM events")
Или переименуйте колонки при создании view:
df.withColumnRenamed("end", "end_date") \
.withColumnRenamed("value", "metric_value") \
.withColumnRenamed("table", "source_table") \
.createOrReplaceTempView("events")
Чеклист¶
Перед включением ANSI mode:
- Запустить аудит всех ETL-запросов с
ansi.enabled=trueв dev среде - Собрать все
ArithmeticException,AnalysisException,CAST_INVALID_INPUT - Проверить все деления - заменить на
try_divideили добавитьCASE WHEN - Проверить все
CAST- заменить наTRY_CASTдля внешних данных - Заменить
IntegerTypeнаLongTypeвезде, где возможен overflow - Добавить явные
CASTвместо implicit coercion - Экранировать зарезервированные слова в именах колонок
После включения:
- Настроить мониторинг
null-rate послеTRY_CAST - Добавить quarantine sink для невалидных записей
- Обновить тесты - добавить ANSI mode в test SparkSession
- Настроить
storeAssignmentPolicy=ANSIдля Delta/Iceberg таблиц
Конфигурации для production:
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.conf.set("spark.sql.storeAssignmentPolicy", "ANSI")
spark.conf.set("spark.sql.datetime.java8API.enabled", "true")
spark.conf.set("spark.sql.session.timeZone", "UTC")
Итог¶
ANSI mode - это не просто переключатель, это философия разработки. Legacy Spark выбирает молчаливость: ошибки поглощаются, данные продолжают течь. ANSI mode выбирает честность: проблема должна быть видна разработчику немедленно.
Главный принцип: silent data corruption хуже падения пайплайна. Пайплайн, который упал, заметен сразу и исправляется. Пайплайн, который тихо выдаёт неверные данные, может годами вводить бизнес в заблуждение.
Функции try_divide, try_cast, try_multiply, try_element_at - это не компромисс с Legacy, а правильный инструментарий для работы с реальными данными: используйте их там, где null - ожидаемый и корректный результат. Используйте строгие cast и прямые операции там, где ошибка - знак о проблеме в данных.
Домашнее задание¶
Задача 1. Возьмите любой существующий PySpark-скрипт и включите spark.sql.ansi.enabled=true. Какие ошибки возникли? Классифицируйте их по типам из урока и предложите исправления.
Задача 2. Реализуйте функцию safe_ratio(numerator_col, denominator_col, default_value=None), которая возвращает колонку с безопасным делением. Функция должна:
- Возвращать
default_valueпри делении на ноль (а неnull) - Работать в ANSI mode
- Принимать DataFrame и имена колонок в качестве аргументов
Задача 3. Напишите Data Quality checker - класс или функцию, которая принимает DataFrame и schema (словарь {column: expected_type}), проверяет конвертацию через try_cast и возвращает:
valid_df- строки, где все поля успешно сконвертированыinvalid_df- строки с хотя бы однимnullпосле конвертацииquality_report- словарь сnull_ratesпо каждой колонке