Azure Blob存储中Python脚本处理患者数据时CSV重复追加历史数据的问题求助
Azure Blob存储中Python脚本处理患者数据时CSV重复追加历史数据的问题求助
看起来你遇到的问题是每次处理新患者数据时,都会把之前所有的历史数据重复追加到Azure Blob的CSV文件里,而不是只添加新的患者记录。我帮你分析下代码里的核心问题,再给出具体的修复方案:
问题根源分析
- Kafka消息处理逻辑错误:你代码里用
consumer.poll(1.0)获取消息后,直接把返回值当作单个消息处理,但实际上poll()返回的是一个字典(键是TopicPartition,值是消息列表),这种错误的处理方式会导致消息遍历混乱,甚至重复读取。 - Offset管理不当:你手动分配了分区,并且设置了
auto.offset.reset: earliest,但没有持久化消费偏移量。这意味着每次启动消费者时,都会从头开始消费所有历史消息,把之前已经处理过的患者数据再次追加到CSV里。
修复方案
方案一:用消费者组自动管理Offset(推荐)
这种方式让Kafka自动跟踪你的消费位置,无需手动维护偏移量,是最省心的做法:
修改validate_and_save_csv_data函数中的Kafka相关代码:
def validate_and_save_csv_data(self, resource, resource_blob): err = [] try: model = apps.get_model( "api", str(resource.get("landing_table_name")).replace("_", "") ) obj_list = model.objects.filter(ready_to_process=True) fhir_id_list = list(set([obj.fhir_id for obj in obj_list if obj.fhir_id])) conf = { "bootstrap.servers": settings.KAFKA_SERVER, "group.id": settings.KAFKA_GROUP_ID, "auto.offset.reset": "earliest", # 首次运行从头消费,后续从上次位置继续 "enable.auto.commit": False, "session.timeout.ms": 10000, # 确保消费者组正常工作的必要配置 } consumer = Consumer(conf) # 替换手动assign为subscribe,让消费者组自动管理offset topic_name = f"cdr_task_{resource.get('resource_name', '')}" consumer.subscribe([topic_name]) data = [] try: msg_count = 0 while True: # poll返回的是{TopicPartition: [ConsumerRecord, ...]}的字典 msg_dict = consumer.poll(timeout_ms=1000) if not msg_dict: # 没有新消息时退出循环 break for tp, records in msg_dict.items(): for record in records: msg_count += 1 if record.error(): if record.error().code() == KafkaError._PARTITION_EOF: logger.debug(f"已读取到分区{tp}的末尾") continue else: raise KafkaException(record.error()) # 处理消息内容 row = record.value().decode("utf-8").split(",") valid = False for val in row: if val and val in fhir_id_list: valid = True break if valid: data.append(row) logger.debug(f"添加有效记录: {row}") # 手动提交当前消息的offset(+1表示下一次从这条消息的下一条开始) consumer.commit({tp: OffsetAndMetadata(record.offset() + 1, None)}) logger.debug(f"已提交分区{tp}的偏移量: {record.offset() + 1}") except Exception as e: err.append( { "severity": "error", "code": "code-invalid", "details": {"text": str(e)}, "diagnostics": "生成CSV数据时发生错误。", } ) finally: consumer.close() if data: append_to_csv_in_blob(resource_blob, data) except Exception as e: logger.error( "处理资源" + str(resource.get("resource_name")) + "的消息时发生错误: " + str(e) ) return err
方案二:手动维护Offset(适合特殊业务场景)
如果必须手动控制消费位置,可以把偏移量存储到Blob的元数据中,每次启动时读取并定位:
- 新增偏移量读写函数:
def save_blob_offset(blob_name, offset): container_client = BlobServiceClient.from_connection_string(AZURE_STORAGE_CONNECTION_STRING) blob_client = container_client.get_blob_client(container=AZURE_BLOB_CONTAINER_NAME, blob=blob_name) # 更新Blob元数据 metadata = blob_client.get_blob_properties().metadata metadata["last_offset"] = str(offset) blob_client.set_blob_metadata(metadata) def get_blob_offset(blob_name): container_client = BlobServiceClient.from_connection_string(AZURE_STORAGE_CONNECTION_STRING) blob_client = container_client.get_blob_client(container=AZURE_BLOB_CONTAINER_NAME, blob=blob_name) try: metadata = blob_client.get_blob_properties().metadata return int(metadata.get("last_offset", 0)) except: return 0
- 修改消费逻辑:
# 在创建consumer后添加以下代码 topic_name = f"cdr_task_{resource.get('resource_name', '')}" topic_partition = TopicPartition(topic_name, 0) consumer.assign([topic_partition]) # 读取上次保存的偏移量并定位 saved_offset = get_blob_offset(resource_blob) consumer.seek(topic_partition, saved_offset) # 在commit offset后添加保存逻辑 consumer.commit({tp: OffsetAndMetadata(record.offset() + 1, None)}) save_blob_offset(resource_blob, record.offset() + 1)
效果说明
修复后,每次处理患者数据时,Kafka消费者只会读取新生成的消息,不会重复消费历史数据,这样append_to_csv_in_blob函数只会把新的患者记录追加到CSV中,就不会出现历史数据重复的问题了。
备注:内容来源于stack exchange,提问作者krishna sai
相关产品推荐
相关产品推荐

