You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Azure Blob存储中Python脚本处理患者数据时CSV重复追加历史数据的问题求助

Azure Blob存储中Python脚本处理患者数据时CSV重复追加历史数据的问题求助

看起来你遇到的问题是每次处理新患者数据时,都会把之前所有的历史数据重复追加到Azure Blob的CSV文件里,而不是只添加新的患者记录。我帮你分析下代码里的核心问题,再给出具体的修复方案:

问题根源分析

  1. Kafka消息处理逻辑错误:你代码里用consumer.poll(1.0)获取消息后,直接把返回值当作单个消息处理,但实际上poll()返回的是一个字典(键是TopicPartition,值是消息列表),这种错误的处理方式会导致消息遍历混乱,甚至重复读取。
  2. 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的元数据中,每次启动时读取并定位:

  1. 新增偏移量读写函数:
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
  1. 修改消费逻辑:
# 在创建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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.13 20:04:34