如何用Python从Tally报表取数并同步至Google BigQuery(含CDC)
问题解决:Tally数据转存BigQuery失败+CDC实现
一、BigQuery写入失败排查
先从最常见的问题入手定位:
- 数据类型不匹配:检查BigQuery表结构和JSON字段类型是否对齐,比如Tally返回的数值是字符串格式,BigQuery字段设为
INT64就会报错,需要提前做类型转换(比如int(record['Amount']))。 - 字段不一致:JSON里有BigQuery表未定义的字段,或者必填字段为空,需要提前过滤多余字段、补全必填项,或者修改BigQuery表结构。
- 权限不足:确认你的服务账号拥有
bigquery.tables.update和bigquery.jobs.create权限,去IAM控制台检查角色配置。 - XML转JSON异常:Tally的XML嵌套结构复杂,转JSON时可能把数组转成单个对象,用
json.dumps()输出转换后的JSON,检查是否有格式错误。 - 批量写入限制:单次写入行数过多会触发BigQuery限制,拆分批次(比如每1000行一批),或者改用流式插入(注意成本)。
二、CDC(变更数据捕获)实现步骤
要实现新增、更新、删除的识别,核心是靠唯一标识和数据对比:
- 确定唯一主键:从Tally数据里找唯一标识,比如凭证号
VoucherNumber、 ledger IDLedgerID,用来匹配BigQuery已有数据。 - 拉取现有数据快照:查询BigQuery目标表的主键和数据特征(比如字段哈希值、最后更新时间)。
- 对比生成变更集:
- 新增:Tally有但BigQuery没有的主键记录
- 更新:主键存在,但字段值有变化(用哈希值对比所有字段,或用Tally返回的
LastModifiedDate筛选) - 删除:BigQuery有但Tally没有的主键记录(注意:Tally默认不返回已删除数据,要么定期全量同步标记删除,要么在Tally开启删除日志)
- 执行变更操作:分别用
INSERT、MERGE、DELETE语句操作BigQuery。
三、示例代码
import xmltodict import json from google.cloud import bigquery import hashlib # 1. Tally XML转JSON def xml_to_json(xml_response): xml_dict = xmltodict.parse(xml_response) # 提取Tally报表核心数据(根据你的XML结构调整路径) tally_records = xml_dict.get('ENVELOPE', {}).get('BODY', {}).get('DATA', {}).get('TALLYMESSAGE', []) # 转成标准JSON格式 return json.loads(json.dumps(tally_records)) # 2. 初始化BigQuery客户端 client = bigquery.Client() table_id = "你的项目ID.数据集ID.表名" # 3. 获取BigQuery现有数据的主键和哈希值 def get_existing_data(): # 替换成你的表字段,生成数据哈希用于对比更新 field_list = ['VoucherNumber', 'LedgerName', 'Amount', 'PostingDate'] hash_sql = f"SHA256(CONCAT({','.join([f'CAST({col} AS STRING)' for col in field_list])}))" query = f"SELECT VoucherNumber, {hash_sql} AS data_hash FROM `{table_id}`" query_job = client.query(query) return {row.VoucherNumber: row.data_hash for row in query_job.result()} # 4. CDC逻辑处理 def process_cdc(tally_data, existing_data): inserts = [] updates = [] deletes = [] tally_ids = {record['VoucherNumber'] for record in tally_data} # 处理新增和更新 for record in tally_data: record_id = record['VoucherNumber'] # 生成当前记录的哈希值 current_hash = hashlib.sha256(json.dumps(record, sort_keys=True).encode()).hexdigest() if record_id not in existing_data: inserts.append(record) else: if existing_data[record_id] != current_hash: updates.append(record) # 处理删除(BigQuery有但Tally没有的记录) delete_ids = existing_data.keys() - tally_ids deletes = list(delete_ids) return inserts, updates, deletes # 5. 写入BigQuery def write_changes(inserts, updates, deletes): # 插入新增数据 if inserts: errors = client.insert_rows_json(table_id, inserts) if errors: print(f"插入失败: {errors}") # 更新数据(用MERGE语句) if updates: update_fields = ', '.join([f'target.{k} = source.{k}' for k in updates[0].keys() if k != 'VoucherNumber']) merge_sql = f""" MERGE `{table_id}` AS target USING UNNEST({json.dumps(updates)}) AS source ON target.VoucherNumber = source.VoucherNumber WHEN MATCHED THEN UPDATE SET {update_fields} """ client.query(merge_sql).result() # 删除数据 if deletes: delete_sql = f""" DELETE FROM `{table_id}` WHERE VoucherNumber IN UNNEST({json.dumps(deletes)}) """ client.query(delete_sql).result() # 主流程 if __name__ == "__main__": # 替换成你的Tally XML响应内容 tally_xml = """<ENVELOPE>...</ENVELOPE>""" tally_json_data = xml_to_json(tally_xml) existing_records = get_existing_data() inserts, updates, deletes = process_cdc(tally_json_data, existing_records) write_changes(inserts, updates, deletes)
四、关键注意事项
- 性能优化:如果Tally返回
LastModifiedDate,可以只同步上次同步时间之后的记录,不用全量对比哈希,提升效率。 - 删除逻辑验证:Tally默认不返回已删除凭证,若需处理删除,要么定期全量同步标记删除状态,要么在Tally系统中开启删除日志功能。
- 错误重试:给写入操作加重试机制(比如用
tenacity库),同时把错误记录到日志表,方便排查。 - 表结构维护:每次Tally报表字段变更时,同步更新BigQuery表结构,避免类型不匹配问题。
内容的提问来源于stack exchange,提问作者Abin Benny
相关产品推荐
相关产品推荐

