Я пытаюсь оптимизировать поток SCD2 в Spark SQL (Python API) с помощью кэшированного промежуточного DataFrame.
Среда
- Python 3.10
- Spark 3.3.2
- Spark SQL (Python)
- Конфигурации Spark (spark.driver.memory = 12g, master = local[6])
- Таблица Iceberg в качестве цели
- Источник — «большой» набор данных (3,2 ГБ)
- Я вижу сканы источника с помощью 12 разделов (так же, как количество исходных разделов) на этапах, где я ожидал повторного использования кэша.
Я создаю промежуточный DataFrame (
Код: Выделить всё
staged_deltaЯ материализую кеш с помощью count().
Затем я использую это кэшированное представление в двух действиях:
- , чтобы закрыть текущие записи (
Код: Выделить всё
MERGE)Код: Выделить всё
delta_type = 'U' - , чтобы открыть новые версии (
Код: Выделить всё
INSERT)Код: Выделить всё
delta_type IN ('I','U')
Упрощенный и анонимизированный код
Код: Выделить всё
from pyspark import StorageLevel
def apply_scd2(spark, staged_df, target_table):
# 1) Prepare source snapshot for this batch
staged_df = staged_df.repartition(500)
staged_df.createOrReplaceTempView("source_staged")
# 2) Current target slice
spark.sql(f"""
CREATE OR REPLACE TEMP VIEW current_target AS
SELECT key_a, key_b, period_key, change_hash
FROM {target_table}
WHERE is_current = true
""")
# 3) Build delta and cache it
staged_delta = spark.sql("""
SELECT
s.*,
t.change_hash AS current_hash,
CASE
WHEN t.key_a IS NULL THEN 'I'
WHEN t.change_hash s.change_hash THEN 'U'
ELSE 'N'
END AS delta_type
FROM source_staged s
LEFT JOIN current_target t
ON t.key_a = s.key_a
AND t.key_b = s.key_b
AND t.period_key = s.period_key
""").persist(StorageLevel.MEMORY_AND_DISK)
# Materialize cache
staged_delta.count()
staged_delta.createOrReplaceTempView("staged_delta")
# 4) Close rows (updates)
spark.sql("""
CREATE OR REPLACE TEMP VIEW close_source AS
SELECT key_a, key_b, period_key, batch_ts, batch_id, delta_type
FROM staged_delta
WHERE delta_type = 'U'
""")
spark.sql(f"""
MERGE INTO {target_table} t
USING close_source s
ON t.key_a = s.key_a
AND t.key_b = s.key_b
AND t.period_key = s.period_key
AND t.is_current = true
WHEN MATCHED AND s.delta_type = 'U' THEN UPDATE SET
t.valid_to = s.batch_ts,
t.is_current = false,
t.ingestion_ts = current_timestamp(),
t.batch_id = s.batch_id
""")
# 5) Insert open rows (inserts)
spark.sql(f"""
INSERT INTO {target_table} (
key_a, key_b, period_key, attr_1, attr_2,
valid_from, valid_to, is_current,
change_hash, ingestion_ts, batch_id
)
SELECT
s.key_a, s.key_b, s.period_key, s.attr_1, s.attr_2,
s.batch_ts, CAST(NULL AS TIMESTAMP), true,
s.change_hash, current_timestamp(), s.batch_id
FROM staged_delta s
WHERE s.delta_type IN ('I', 'U')
""")
- кэшируется один раз.
Код: Выделить всё
staged_delta - и INSERT должны использовать кэшированный staged_delta без дорогостоящего повторного сканирования и перемешивания источника.
Код: Выделить всё
MERGE
- Во время MERGE и/или INSERT Spark запускает сканирование, аналогичное шаблону сканирования источника.
- Я все еще вижу этапы с количеством разделов, соответствующим разделению источника. (12), предлагающее перерасчет или возврат к происхождению.
Почему Spark не использует кэшированный DataFrame?
Примечания
[*]Действительные имена бизнес-столбцов анонимизируются.
[*]При необходимости я могу предоставить физические планы.