Databricks DLT DataFrame:视图场景下如何通过Schema设置列注释
问题描述
我是Databricks Delta Live Tables(DLT)和DataFrame的新手,在表到表的流处理场景中,对读取流时Schema的使用存在困惑。我们需要让输出表的列注释能在DBX Catalog概览页面正常显示。
输入表的简化Schema如下:
message_schema = StructType([ StructField("bayId", StringType(), True, {"comment": "Bay ID"}), StructField("storeId", StringType(), True, {"comment": "Store ID"}), StructField("message", StringType(), True, {"comment": "Message content"}), ])
其中message列包含JSON字符串,对应的简化Schema为:
event_schema = StructType([ StructField("Id", StringType(), True, {"comment": "Event ID"}), StructField("Payload", StructType([ StructField("PlayerDexterity", StringType(), True, {"comment": "Player dexterity"}), StructField("AttackAngle", FloatType(), True, {"comment": "The vertical angle at which the club head approaches the ball"}), ]), ])
我的流处理代码如下:
df = spark.readStream.table("tablename") df = df.where( col("MESSAGE").isNotNull() ).select( col("BAYID"), col("STOREID"), F.from_json(col("MESSAGE"), message_schema).alias("MESSAGE") ).select( col("BAYID"), col("STOREID"), F.from_json(col("MESSAGE.message"), event_schema).alias("EVENT") ).select( F.expr("uuid()").alias("ID"), col("BAYID").alias("BAY_ID"), col("STOREID").alias("STORE_ID"), col("EVENT.Payload.PlayerDexterity").alias("PLAYER_DEXTERITY_NAME"), col("EVENT.Payload.AttackAngle").alias("ATTACK_ANGLE_NBR") )
运行后发现,输出表在Catalog概览页中,PLAYER_DEXTERITY_NAME和ATTACK_ANGLE_NBR列能显示Schema中设置的注释,但BAY_ID和STORE_ID列没有注释。我通过添加以下代码成功给这两列加上了注释:
df = (df .withMetadata("BAY_ID", {"comment": "Bay ID"}) .withMetadata("STORE_ID", {"comment": "Store ID"}) )
但为了保持一致性,我希望直接通过Schema来设置注释。尝试在spark.readStream中指定schema(message_schema)没有效果,而且因为需要使用@dlt.view(),无法在@dlt.table()中指定Schema,请问该如何解决?
解决方案
方法1:定义最终输出Schema并强制应用
先创建包含所有输出列及对应注释的最终Schema,再将处理后的DataFrame转换为该Schema,以此保留注释元数据:
- 定义最终输出Schema:
final_schema = StructType([ StructField("ID", StringType(), True, {"comment": "Auto-generated UUID"}), StructField("BAY_ID", StringType(), True, {"comment": "Bay ID"}), StructField("STORE_ID", StringType(), True, {"comment": "Store ID"}), StructField("PLAYER_DEXTERITY_NAME", StringType(), True, {"comment": "Player dexterity"}), StructField("ATTACK_ANGLE_NBR", FloatType(), True, {"comment": "The vertical angle at which the club head approaches the ball"}), ])
- 在流处理代码的最后一步,用该Schema重新生成DataFrame:
# 替换原有select后的处理逻辑 df = spark.createDataFrame(df.rdd, final_schema)
方法2:直接在DLT表装饰器中声明列注释
利用@dlt.table()的column_comments参数,直接在表定义中指定各列的注释,无需修改DataFrame的Schema逻辑:
@dlt.table( comment="Processed event data with player metrics", column_comments={ "ID": "Auto-generated UUID", "BAY_ID": "Bay ID", "STORE_ID": "Store ID", "PLAYER_DEXTERITY_NAME": "Player dexterity", "ATTACK_ANGLE_NBR": "The vertical angle at which the club head approaches the ball" } ) def processed_events(): df = spark.readStream.table("tablename") # 此处保留原有的流处理逻辑 df = df.where( col("MESSAGE").isNotNull() ).select( col("BAYID"), col("STOREID"), F.from_json(col("MESSAGE"), message_schema).alias("MESSAGE") ).select( col("BAYID"), col("STOREID"), F.from_json(col("MESSAGE.message"), event_schema).alias("EVENT") ).select( F.expr("uuid()").alias("ID"), col("BAYID").alias("BAY_ID"), col("STOREID").alias("STORE_ID"), col("EVENT.Payload.PlayerDexterity").alias("PLAYER_DEXTERITY_NAME"), col("EVENT.Payload.AttackAngle").alias("ATTACK_ANGLE_NBR") ) return df
方法3:从输入表继承列注释元数据
直接读取输入表中BAYID和STOREID的注释,在重命名列时传递这些元数据,避免重复编写注释内容:
def processed_events(): # 读取输入表并获取列注释 input_df = spark.readStream.table("tablename") bay_comment = input_df.schema["BAYID"].metadata.get("comment", "") store_comment = input_df.schema["STOREID"].metadata.get("comment", "") df = input_df.where( col("MESSAGE").isNotNull() ).select( col("BAYID"), col("STOREID"), F.from_json(col("MESSAGE"), message_schema).alias("MESSAGE") ).select( col("BAYID"), col("STOREID"), F.from_json(col("MESSAGE.message"), event_schema).alias("EVENT") ).select( F.expr("uuid()").alias("ID"), # 重命名时传递注释元数据 col("BAYID").alias("BAY_ID", metadata={"comment": bay_comment}), col("STOREID").alias("STORE_ID", metadata={"comment": store_comment}), col("EVENT.Payload.PlayerDexterity").alias("PLAYER_DEXTERITY_NAME"), col("EVENT.Payload.AttackAngle").alias("ATTACK_ANGLE_NBR") ) return df
内容的提问来源于stack exchange,提问作者Westy
相关产品推荐
相关产品推荐

