Есть полезная привычка для ETL и трансформаций: явно приводить типы и сразу фиксировать схему.
😅 Недавно словил неприятную ошибку. Логов было мало, по ним быстро не понять, где именно всё поехало. В итоге понял довольно простой источник проблемы: типы в одном месте получились не те, которые ETL pipeline ожидал увидеть на выходе, и дальше это уже вылезло на вставке.
После этого кейса в таких местах стараюсь не оставлять типы “как получится”. Если на выходе колонка должна быть строкой, я это сразу пишу явно:
CAST(column AS STRING) AS column
Если колонка пока пустая, но дальше она должна лечь как STRING, тоже задаю это сразу:CAST(NULL AS STRING) AS column2
Это особенно помогает, когда собираешь слой под insert, union, промежуточную витрину или просто хочешь получить стабильную схему на выходе.В
Spark та же история. Если читается csv или другой сырой источник, спокойнее сразу задать схему через schema, чтобы потом не разбирать, почему одна колонка внезапно стала string, другая double, а третья вообще прочиталась криво.Например:
from pyspark.sql.types import StructType, StructField, StringType, DecimalType, TimestampType
schema = StructType([
StructField("user_id", StringType(), True),
StructField("amount", DecimalType(38, 2), True),
StructField("event_dt", TimestampType(), True)
])
df = (
spark.read
.option("header", True)
.schema(schema)
.csv("/path/to/file.csv")
)
И потом уже в самой трансформации можно сразу собрать нужный выходной слой:
from pyspark.sql import functions as F
df_final = df.select(
F.col("user_id").cast("string").alias("user_id"), F.col("amount").cast("decimal(38,2)").alias("amount"), F.col("event_dt").cast("timestamp").alias("event_dt"), F.lit(None).cast("decimal(38,2)").alias("bonus_amount")
)
✅ Плюс тут ну прям оооочень жирный. Схема становится предсказуемой, вставки проходят спокойнее,
union не начинает жить своей жизнью, и в самом коде сразу видно, какой слой вы хотите получить на выходе.✍️ Для junior DE это, кстати, хорошая привычка с самого начала. Чем раньше начинаешь смотреть на типы как на часть расчёта, тем меньше потом времени уходит на странные ошибки в пайплайне. Просто поверьте.
💭 Полезные ссылки по теме:
• Spark DataFrameReader.csv
• Spark CSV Files
• PySpark Column.cast
• PySpark StructType
✏️ Резюмируя, если знаешь, какой тип должен быть на выходе, лучше задать его сразу.
#база_знаний