Google Cloud Function实现BigQuery表添加创建日期字段问题
问题描述
我正在用Google Cloud Function创建BigQuery表,当前功能正常,但想给表加一个记录创建时间的字段,目前尝试的代码没生效。核心目标是先实现单表的该功能,再把流程复用到处理两张及以上表的场景中。
当前代码示例
main.py 代码(已修正缩进与语法错误)
from google.cloud import bigquery import pandas as pd from previsional_tables import table_TEST1 # 注意:此处定义的creation_date只会在函数部署时执行一次,无法实现每条记录的创建时间注入 # creation_date = pd.Timestamp.now() def main_function(event, context): dataset = 'bd_clients' file = event input_bucket_name = file['bucket'] path_file = file['name'] uri = 'gs://{}/{}'.format(input_bucket_name, path_file) path_file_list = path_file.split("/") file_name_ext = path_file_list[len(path_file_list) - 1] file_name_ext_list = file_name_ext.split(".") name_file = file_name_ext_list[0] print('nombre archivo ==> ' + name_file.upper()) print('Getting the data from bucket "{}"'.format(uri)) path_file_name = str(uri) print("ruta: ", path_file_name) if ("gs://bucket_test" in path_file_name): client = bigquery.Client() job_config = bigquery.LoadJobConfig() table_test1(dataset, client, uri, job_config, bigquery)
previsional_tables.py 代码(已修正缩进)
def table_test1(dataset, client, uri, job_config, bigquery): table = "test1" dataset_ref = client.dataset(dataset) job_config.autodetect = True job_config.max_bad_records = 1000 job_config.schema = [ bigquery.SchemaField("NAME", "STRING"), bigquery.SchemaField("LAST_NAME", "STRING"), bigquery.SchemaField("ADDRESS", "STRING"), bigquery.SchemaField("DATE", bigquery.enums.SqlTypeNames.DATE) # 定义BigQuery列及类型 ] job_config.source_format = bigquery.SourceFormat.CSV job_config.field_delimiter = ';' job_config.write_disposition = bigquery.WriteDisposition.WRITE_APPEND load_job = client.load_table_from_uri(uri, dataset_ref.table(table), job_config=job_config)
requirements.txt 内容
# Function dependencies, for example: # package>=version google-cloud-bigquery==2.25.1 pysftp==0.2.9 pandas==1.4.2 fsspec==2022.5.0 gcsfs==2022.5.0
相关结构截图
- 存储桶数据结构:

- 函数配置:

- 数据库输出结构:

解决方案
单表添加创建时间字段的实现
你当前的核心问题是:creation_date定义在函数外部(仅部署时执行一次),且未将时间字段注入到加载的数据集里。以下是两种可靠实现方式:
方式1:BigQuery加载时自动注入(推荐,适合大文件)
通过Load Job的Schema更新选项添加字段,再执行SQL补全创建时间:
def table_test1(dataset, client, uri, job_config, bigquery): table = "test1" table_ref = client.dataset(dataset).table(table) # 1. 定义包含创建时间的完整Schema full_schema = [ bigquery.SchemaField("NAME", "STRING"), bigquery.SchemaField("LAST_NAME", "STRING"), bigquery.SchemaField("ADDRESS", "STRING"), bigquery.SchemaField("DATE", "DATE"), bigquery.SchemaField("CREATED_AT", "TIMESTAMP") # 添加创建时间字段 ] # 2. 配置加载规则,允许自动添加字段 job_config.autodetect = True job_config.max_bad_records = 1000 job_config.schema = full_schema job_config.source_format = bigquery.SourceFormat.CSV job_config.field_delimiter = ';' job_config.write_disposition = bigquery.WriteDisposition.WRITE_APPEND job_config.schema_update_options = [ bigquery.SchemaUpdateOption.ALLOW_FIELD_ADDITION ] # 3. 执行数据加载 load_job = client.load_table_from_uri(uri, table_ref, job_config=job_config) load_job.result() # 等待加载完成 # 4. 给刚加载的空值记录补全当前时间 update_sql = f""" UPDATE `{dataset}.{table}` SET CREATED_AT = CURRENT_TIMESTAMP() WHERE CREATED_AT IS NULL """ client.query(update_sql).result()
方式2:Pandas预处理数据(适合小文件)
先将CSV加载到内存,添加时间字段后再上传:
def main_function(event, context): # ... 原有逻辑不变,直到获取uri后 if ("gs://bucket_test" in path_file_name): client = bigquery.Client() # 读取CSV并添加UTC时区的创建时间 df = pd.read_csv(uri, delimiter=';') df['CREATED_AT'] = pd.Timestamp.now(tz='UTC') # 配置上传规则 job_config = bigquery.LoadJobConfig( write_disposition=bigquery.WriteDisposition.WRITE_APPEND, schema=[ bigquery.SchemaField("NAME", "STRING"), bigquery.SchemaField("LAST_NAME", "STRING"), bigquery.SchemaField("ADDRESS", "STRING"), bigquery.SchemaField("DATE", "DATE"), bigquery.SchemaField("CREATED_AT", "TIMESTAMP") ] ) # 上传到BigQuery table_ref = client.dataset('bd_clients').table('test1') load_job = client.load_table_from_dataframe(df, table_ref, job_config=job_config) load_job.result()
多表复用的实现
将创建时间注入逻辑封装为通用函数,根据文件名匹配不同表的Schema:
通用处理函数
def get_full_schema(base_schema): # 给基础Schema添加创建时间字段 base_schema.append(bigquery.SchemaField("CREATED_AT", "TIMESTAMP")) return base_schema def load_table_with_created_at(client, dataset, table_name, uri, base_schema, delimiter=';'): table_ref = client.dataset(dataset).table(table_name) job_config = bigquery.LoadJobConfig( autodetect=True, max_bad_records=1000, schema=get_full_schema(base_schema), source_format=bigquery.SourceFormat.CSV, field_delimiter=delimiter, write_disposition=bigquery.WriteDisposition.WRITE_APPEND, schema_update_options=[bigquery.SchemaUpdateOption.ALLOW_FIELD_ADDITION] ) # 加载数据 load_job = client.load_table_from_uri(uri, table_ref, job_config=job_config) load_job.result() # 补全创建时间 update_sql = f""" UPDATE `{dataset}.{table_name}` SET CREATED_AT = CURRENT_TIMESTAMP() WHERE CREATED_AT IS NULL """ client.query(update_sql).result()
多表调用示例
def main_function(event, context): # ... 原有逻辑不变,直到获取name_file后 if ("gs://bucket_test" in path_file_name): client = bigquery.Client() dataset = 'bd_clients' # 根据文件名匹配不同表的Schema if name_file.upper() == 'TEST1': base_schema = [ bigquery.SchemaField("NAME", "STRING"), bigquery.SchemaField("LAST_NAME", "STRING"), bigquery.SchemaField("ADDRESS", "STRING"), bigquery.SchemaField("DATE", "DATE") ] load_table_with_created_at(client, dataset, 'test1', uri, base_schema) elif name_file.upper() == 'TEST2': base_schema = [ bigquery.SchemaField("ID", "INTEGER"), bigquery.SchemaField("PRODUCT", "STRING"), bigquery.SchemaField("PRICE", "FLOAT") ] load_table_with_created_at(client, dataset, 'test2', uri, base_schema)
内容的提问来源于stack exchange,提问作者Nelson Romero
相关产品推荐
相关产品推荐

