Spark/Databricks中基于条件生成列并添加规则记录列的问题
问题解决:Spark动态生成Rules_applied列时的TypeError错误
错误原因
你遇到的TypeError: sequence item 0: expected str instance, Column found,是因为在expr("concat_ws(',', {})".format(...))中,尝试将when(col(rule_expr), lit(rule_name))生成的Column对象直接拼接成字符串,但Column不是字符串类型,无法被str.join()处理。
另外注意代码里的笔误:validity_rules.values()应该改为rules.values(),因为你定义的规则字典是rules。
解决方案
方案1:使用Spark Column API构建(推荐)
直接用Spark的Column操作生成规则列,避免字符串拼接的问题:
from pyspark.sql.functions import concat_ws, when, lit, expr rules = { "valid_website": "(Website IS NOT NULL)", "valid_location": "(Location IS NOT NULL)" } quarantine_rules = "NOT({0})".format(" AND ".join(rules.values())) @dlt.table(name="valid_application") @dlt.expect_all(rules) def get_quarantined_application(): # 为每个规则生成对应Column:满足规则返回规则名,否则返回null rule_columns = [when(expr(rule_expr), lit(rule_name)) for rule_name, rule_expr in rules.items()] return ( dlt.read("msf_raw_application") .withColumn("is_valid", expr(quarantine_rules)) # 用concat_ws拼接非null的规则名,自动忽略null值 .withColumn("Rules_applied", concat_ws(",", *rule_columns)) )
方案2:拼接字符串表达式传入expr
如果坚持用expr方式,需要把每个规则的判断逻辑转成字符串形式的SQL表达式:
from pyspark.sql.functions import expr rules = { "valid_website": "(Website IS NOT NULL)", "valid_location": "(Location IS NOT NULL)" } quarantine_rules = "NOT({0})".format(" AND ".join(rules.values())) # 生成每个规则的SQL IF表达式字符串 rule_expr_strings = [f"IF({rule_expr}, '{rule_name}', NULL)" for rule_name, rule_expr in rules.items()] # 拼接成concat_ws的完整SQL表达式 concat_sql = f"concat_ws(',', {', '.join(rule_expr_strings)})" @dlt.table(name="valid_application") @dlt.expect_all(rules) def get_quarantined_application(): return ( dlt.read("msf_raw_application") .withColumn("is_valid", expr(quarantine_rules)) .withColumn("Rules_applied", expr(concat_sql)) )
两种方案最终都会生成Rules_applied列,记录所有触发的规则名称(用逗号分隔),比如当Website和Location都非空时,该列值为valid_website,valid_location;如果只有Website非空,则值为valid_website。
内容的提问来源于stack exchange,提问作者wingrider
相关产品推荐
相关产品推荐

