如何在Spark中修复因类型转换导致的无效输入错误?
我在Databricks中尝试创建并显示一个Spark DataFrame。这是我的代码的样子:
df = spark.sql(f'''
SELECT A.CUSTOMER_ID
FROM TABLE_1 A
INNER JOIN
TABLE_2 B
ON A.PRODUCT_ID = B.PRODUCT_ID
WHERE A.ORDER_DATE BETWEEN DATE_SUB(CURRENT_DATE, 182) AND CURRENT_DATE
AND A.CUSTOMER_ID IS NOT NULL
AND B.ITEM_STATUS IN ('Live', 'Active')
AND B.ITEM_CODE IN (2009, 2012)
GROUP BY A.CUSTOMER_ID
''')
df.display()
在按上述方式运行时会返回以下错误:
[CAST_INVALID_INPUT] 值为 'UNKNOWN' 且类型为 "STRING" 的值无法转换为 "BIGINT",因为它的格式不正确。
我理解这意味着Spark正在对某个字符串列执行隐式转换,但它失败了,因为该列包含值 'UNKNOWN',无法转换为bigint的格式。
所以我尝试从我的查询使用的每个字符串列中移除值 'UNKNOWN':
df = spark.sql(f'''
SELECT A.CUSTOMER_ID
FROM (
SELECT *
FROM TABLE_1
WHERE UPPER(PRODUCT_ID) <> 'UNKNOWN'
AND UPPER(ORDER_DATE) <> 'UNKNOWN'
AND UPPER(CUSTOMER_ID) <> 'UNKNOWN'
) A
INNER JOIN
(
SELECT *
FROM TABLE_2
WHERE UPPER(PRODUCT_ID) <> 'UNKNOWN'
AND UPPER(ITEM_STATUS) <> 'UNKNOWN'
AND UPPER(ITEM_CODE) <> 'UNKNOWN'
) B
ON A.PRODUCT_ID = B.PRODUCT_ID
WHERE A.ORDER_DATE BETWEEN DATE_SUB(CURRENT_DATE, 182) AND CURRENT_DATE
AND A.CUSTOMER_ID IS NOT NULL
AND B.ITEM_STATUS IN ('Live', 'Active')
AND B.ITEM_CODE IN (2009, 2012)
GROUP BY A.CUSTOMER_ID
''')
df.display()
然而,我仍然得到同样的错误信息。我不明白这是怎么回事,因为未知值已经被移除。谁能帮助解释到底发生了什么,以及如何解决这个错误?
解决方案
你可以通过在 ITEM_CODE 的比较中使用字符串值来替代整数,从而解决这个问题:
B.ITEM_CODE IN ('2009', '2012')
这可以阻止Spark将 ITEM_CODE 列隐式转换为 BIGINT。
至于为什么在子查询中过滤掉 'UNKNOWN' 行后你的第二个查询仍然失败——这是因为Spark的 Catalyst优化器会构建自己的查询执行计划。它不一定先执行子查询的过滤条件再应用外部WHERE子句。
优化器可以自由地合并、重新排序以及向下推送谓词,因此对ITEM_CODE的隐式转换仍可能在包含 'UNKNOWN' 的行上尝试,甚至在这些行被过滤掉之前。
如果你仍然需要将 ITEM_CODE 作为整数进行比较,可以使用 TRY_CAST。TRY_CAST 返回 NULL,而不是在遇到无法转换的值时抛出错误。
AND TRY_CAST(B.ITEM_CODE AS BIGINT) IN (2009, 2012)
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。