Azure SQL CDC列含特殊字符导致PySpark Delta MERGE执行失败
解决Azure SQL CDC到Delta Lake MERGE语句的特殊列名语法错误
Spark SQL处理含$这类特殊字符的列名时,必须用**反引号(`)**转义——单引号或双引号仅用于包裹字符串常量,无法转义列名。针对你的场景,修改MERGE语句如下:
MERGE INTO poc_table_2 AS target USING cdc_data AS source ON target.last_processed_lsn = source.`__$start_lsn` WHEN MATCHED AND source.`__$operation` = 1 THEN DELETE WHEN MATCHED AND source.`__$operation` IN (3, 4) THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *
额外补充(DataFrame API方案)
如果不想写SQL字符串,也可以用Delta Lake的DataFrame API实现相同逻辑,避免语法转义问题:
from delta.tables import DeltaTable delta_table = DeltaTable.forPath(spark, "/path/to/poc_table_2") delta_table.alias("target").merge( source=cdc_data.alias("source"), condition="target.last_processed_lsn = source.`__$start_lsn`" ).whenMatchedDelete(condition="source.`__$operation` = 1")\ .whenMatchedUpdateAll(condition="source.`__$operation` IN (3, 4)")\ .whenNotMatchedInsertAll()\ .execute()
内容的提问来源于stack exchange,提问作者Eugene Goldberg
相关产品推荐
相关产品推荐

