Spark SQL MERGE/INSERT на Iceberg пересчитывает соединение восходящего потока вместо повторного использования кэшированнPython

Программы на Python
Anonymous
Spark SQL MERGE/INSERT на Iceberg пересчитывает соединение восходящего потока вместо повторного использования кэшированн

Сообщение Anonymous »

Spark SQL + Iceberg: MERGE и INSERT игнорируют кэшированный DataFrame и повторно сканируют источник
Я пытаюсь оптимизировать поток 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
) из соединения и сохраняю его с помощью MEMORY_AND_DISK.
Я материализую кеш с помощью count().
Затем я использую это кэшированное представление в двух действиях: Во время выполнения похоже, что Spark снова сканирует исходные данные на этапах слияния/вставки вместо полного повторного использования кэшированного staged_delta.
Упрощенный и анонимизированный код

Код: Выделить всё

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 и/или INSERT Spark запускает сканирование, аналогичное шаблону сканирования источника.
  • Я все еще вижу этапы с количеством разделов, соответствующим разделению источника. (12), предлагающее перерасчет или возврат к происхождению.
В чем мне нужна помощь
Почему Spark не использует кэшированный DataFrame?
Примечания

[*]Действительные имена бизнес-столбцов анонимизируются.
[*]При необходимости я могу предоставить физические планы.

Вернуться в «Python»