使用PySpark BigTable连接器时,如何读写单元格的时间戳?
关于Spark-BigTable连接器读写单元格时间戳的问题
我正在使用spark-bigtable_2.12-0.6.0.jar连接器,通过PySpark读取BigTable中的数据,使用的示例代码如下:
catalog = f""" {{ "table": {{ "name": "{table_id}", "tableCoder": "PrimitiveType" }}, "rowkey": "id_rowkey", "columns": {{ "id_rowkey": {{ "cf": "rowkey", "col": "id_rowkey", "type": "string" }}, "date_{column}": {{ "cf": "{column}", "col": "key", "type": "string" }}, "{column}": {{ "cf": "{column}", "col": "value", "type": "binary" }}, }} }}""" readDf = spark.read \ .format('bigtable') \ .option('spark.bigtable.project.id', bt_project_id) \ .option('spark.bigtable.instance.id', bt_instance_id) \ .options(catalog=catalog) \ .load()
目前已成功读取数据,但需要访问每个单元格的值与时间戳。已知该连接器支持通过spark.bigtable.read.timerange.start.milliseconds过滤特定时间范围,以及通过spark.bigtable.write.timestamp.milliseconds写入指定时间戳,请问是否有方法通过该连接器读写每个单元格的时间戳?
解决方案
1. 读取单元格时间戳
在catalog配置中,给需要获取时间戳的列添加"includeTimestamp": true参数,连接器会自动为该列生成一个后缀为_timestamp的字段,类型为long(存储毫秒级时间戳)。
修改后的catalog示例:
catalog = f""" {{ "table": {{ "name": "{table_id}", "tableCoder": "PrimitiveType" }}, "rowkey": "id_rowkey", "columns": {{ "id_rowkey": {{ "cf": "rowkey", "col": "id_rowkey", "type": "string" }}, "date_{column}": {{ "cf": "{column}", "col": "key", "type": "string", "includeTimestamp": true }}, "{column}": {{ "cf": "{column}", "col": "value", "type": "binary", "includeTimestamp": true }}, }} }}"""
加载后,DataFrame会包含对应时间戳字段:
date_{column}_timestamp:date_{column}列单元格的毫秒级时间戳{column}_timestamp:{column}列单元格的毫秒级时间戳
直接通过这些_timestamp字段即可访问每个单元格的时间戳。
2. 写入单元格时间戳
如果需要为每个单元格指定不同的时间戳,只需在待写入的DataFrame中添加对应列的{列名}_timestamp字段(类型为long,毫秒级时间戳),连接器会自动将该时间戳应用到对应的单元格上,无需额外配置全局时间戳参数。
示例写入代码:
# 假设DataFrame包含值字段和对应的时间戳字段 writeDf = readDf.select("id_rowkey", "date_{column}", "date_{column}_timestamp", "{column}", "{column}_timestamp") writeDf.write \ .format('bigtable') \ .option('spark.bigtable.project.id', bt_project_id) \ .option('spark.bigtable.instance.id', bt_instance_id) \ .options(catalog=catalog) \ .mode("overwrite") \ .save()
注意:如果同时配置了全局的spark.bigtable.write.timestamp.milliseconds参数,DataFrame中存在的_timestamp字段会覆盖全局配置的时间戳。
内容的提问来源于stack exchange,提问作者El Mahdi KHACHAFI
相关产品推荐
相关产品推荐

