OpenMed 分布式 SQL UDF:在分布式 SQL 引擎内以向量化方式完成临床文本去标识化
【免费下载链接】openmedLocal-first healthcare AI: clinical NER & HIPAA PII de-identification that runs 100% on-device. 2,200+ medical models, 21 languages, Apple MLX + Python, no cloud, no patient data leaving your network. Apache-2.0项目地址: https://gitcode.com/GitHub_Trending/ope/openmed
导读
本文讲解 OpenMed 提供的分布式 SQL UDF 接入层(distributed SQL UDF bridge):通过一个引擎无关的 Python 入口,让 Spark、Ray、Dataflow 等分布式 SQL 引擎的 Python worker 进程直接调用 OpenMed 的去标识化能力。读者将掌握openmed_deidentify(text, profile) -> VARCHAR标量 SQL 函数背后的注册描述符、worker 生命周期、向量化批处理语义(NULL/空串透传、行序保持、profile 分组),以及如何基于源码与测试验证这一接入层的正确性与 PHI 安全边界。
设计目标:把"每行一次模型初始化"变成"每个 worker 一份常驻模型"
OpenMed 以本地优先(local-first)为核心:临床 NER 与 HIPAA PII 去标识化 100% 在设备/集群内运行,患者数据不离开网络。当去标识化流水线需要嵌入分布式 SQL 引擎(数据湖、数仓、流处理平台)时,最差的实现方式是"每个 SQL 行初始化一次模型、做一次推理"——模型加载开销会完全淹没推理开销。
OpenMed 的分布式 SQL UDF 接入层通过一个逻辑上标量、执行上向量化的 SQL 函数规避这一问题:
openmed_deidentify(text VARCHAR, profile VARCHAR) -> VARCHAR该 SQL 函数逻辑上是标量的,但 Python bridge 必须按**向量模式(vectorized mode)调用底层实现。每个 worker 进程只持有一个懒加载(lazy)**的 OpenMed 模型加载器,行窗口(row window)通过process_batch成批送过模型,从而把"N 次模型初始化 + N 次推理调用"压缩为"1 次初始化 + ceil(N/batch_size) 次推理调用"。源码 distributed_sql_udf.py 顶部注释对此做了明确声明。
需要特别说明的是:OpenMed 不随包分发任何编译好的引擎插件 JAR。引擎专属 Python bridge 的供应(provisioning)与安全加固(securing)是运维方的责任,OpenMed 提供的是引擎无关的 Python 入口与注册描述符,由运维方将其映射到目标引擎的 Python UDF 注册体系。
注册描述符:OPENMED_DEIDENTIFY_DESCRIPTOR
引擎无关的注册描述符以模块级常量OPENMED_DEIDENTIFY_DESCRIPTOR形式暴露(定义于 distributed_sql_udf.py):
from openmed.integrations.distributed_sql_udf import ( OPENMED_DEIDENTIFY_DESCRIPTOR, ) descriptor = OPENMED_DEIDENTIFY_DESCRIPTOR其注册字段如下(与源码中字典的键值一一对应):
name: openmed_deidentify language: python entrypoint: openmed.integrations.distributed_sql_udf:deidentify_batch arguments: - name: text sql_type: VARCHAR python_batch: texts - name: profile sql_type: VARCHAR python_batch: profiles return_type: VARCHAR vectorized: true null_handling: called_on_null_input default_batch_size: 64各字段含义与映射要点:
| 字段 | 值 | 说明 |
|---|---|---|
name | openmed_deidentify | 在 SQL 中可见的函数名 |
language | python | 引擎侧应把该函数注册为 Python UDF |
entrypoint | openmed.integrations.distributed_sql_udf:deidentify_batch | 向量化入口:module:function形式,引擎 worker 进程按此导入 |
arguments | text/profile | 两个VARCHAR参数;python_batch字段声明了引擎应把该列的整列数组以texts/profiles关键字传给入口 |
return_type | VARCHAR | 返回值仍为字符串,与输入行一一对齐 |
vectorized | true | 引擎必须按向量/批量模式调用,而非逐行调用 |
null_handling | called_on_null_input | 有 NULL 输入时函数仍会被调用(由 bridge 内部处理 NULL),引擎不应直接短路为 NULL |
default_batch_size | 64 | 引擎无显式配置时建议使用的批大小,与源码常量DEFAULT_DISTRIBUTED_SQL_BATCH_SIZE一致 |
在把描述符映射到具体引擎(如 Spark 的 pandas UDF、Ray Data 的map_batches、Dataflow 的 DoFn)时,核心要求是:bridge 收到的是text与profile两列的数组,返回数组与输入行一一对应。三个边界语义必须保持:
- SQL
NULL输入 → 输出NULL(不触发模型加载); - 空字符串输入 → 输出空字符串(同样不触发模型加载);
- 非空字符串 → 按 profile 去标识化后返回。
Worker 生命周期:一次构造、全进程复用
对于直接集成或离线注册测试,建议在 worker 进程启动阶段构造一次DistributedSQLDeidentifyUDF可调用对象,之后每个向量窗口都复用它:
from openmed.integrations.distributed_sql_udf import ( DistributedSQLDeidentifyUDF, DistributedSQLUDFConfig, ) openmed_deidentify = DistributedSQLDeidentifyUDF( config=DistributedSQLUDFConfig( default_profile="hipaa_safe_harbor", batch_size=64, ) ) redacted = openmed_deidentify( [ "Patient Jane Roe has hypertension.", None, "", ], ["hipaa_safe_harbor", None, "hipaa_safe_harbor"], )返回值为["[PERSON] has hypertension.", None, ""]:非空行被去标识化,None保持None,空串保持空串。
DistributedSQLUDFConfig完整配置项
配置类定义于 distributed_sql_udf.py,为不可变(frozen)dataclass,__post_init__会做完整校验:
| 配置项 | 默认值 | 说明 |
|---|---|---|
model_name | OpenMed/OpenMed-PII-SuperClinical-Small-44M-v1 | 模型注册表键或 Hugging Face 标识符;离线部署时可指向本地缓存的模型构件 |
default_profile | hipaa_safe_harbor | 行未显式指定 profile 时使用的策略名;会经_canonical_profile规范化 |
batch_size | 64 | 单个process_batch窗口的大小,经validate_batch_size校验 |
method | mask | 去标识化方法(如 mask),非空字符串校验 |
confidence_threshold | 0.7 | 触发脱敏的最低置信度,必须落在[0.0, 1.0],否则抛ValueError |
process_batch_kwargs | {} | 透传给底层process_batch的额外关键字参数 |
保留参数:哪些process_batch_kwargs不允许覆盖
process_batch_kwargs用于把引擎特有的额外参数透传给openmed.processing.process_batch。但源码中定义了一个保留集合_RESERVED_PROCESS_BATCH_KWARGS(见 distributed_sql_udf.py),配置校验时若发现这些键出现在process_batch_kwargs中会直接抛ValueError:
batch_size, continue_on_error, loader, method, model_name, operation, policy理由很直接:这些参数由 UDF 层基于DistributedSQLUDFConfig统一注入(operation="deidentify"、continue_on_error=False、policy=profile、method、confidence_threshold、loader),不允许被透传参数静默覆盖,以免破坏行序对齐与错误处理语义。
两种调用形态
deidentify_batch(texts, profiles):向量入口,引擎 bridge 的默认调用点;返回list[str | None],长度与输入行数严格一致。deidentify(text, profile):标量兼容入口,供只能逐行调用的 bridge 使用;其实现本质是deidentify_batch([text], profile)[0]。文档与源码都明确指出:逐行调用无法享受向量化批处理收益,应尽量避免。
模块级函数deidentify与deidentify_batch自动复用同一个进程内 worker——源码通过@lru_cache(maxsize=1)装饰的_default_worker()返回进程局部单例(distributed_sql_udf.py),因此注册描述符中的entrypoint指向deidentify_batch即可,无需在每个引擎事件里重建 worker。
向量化语义的源码级验证
DistributedSQLDeidentifyUDF.__call__与deidentify_batch的处理流程(distributed_sql_udf.py)值得逐段拆解:
- 输入校验:
texts必须是序列(字符串/字节会被拒绝并抛TypeError);profiles可以是单个字符串(广播到所有行)、None或与texts等长的序列——长度不匹配抛ValueError。 - 窗口切分:按
config.batch_size把整列切成多个窗口。 - 窗口内分类:
None→ 输出None并跳过;空字符串 → 输出""并跳过;非空字符串按 profile 归组(grouped_positions字典)。 - 分组推理:同一 profile 的行合并成一次
process_batch调用;返回结果按原始位置写回输出数组,保证行序与输入严格一致。 - 结果校验:
process_batch返回的items数量必须与输入行数一致,否则抛RuntimeError;任何item.error非空也会抛RuntimeError(错误消息刻意不包含源文本)。
profile 名称在进入推理前会经_canonical_profile规范化:转小写、-替换为_,并解析内置别名——_PROFILE_ALIASES中"hipaa"与"safe_harbor"都映射到hipaa_safe_harbor;随后交给openmed.core.policy.canonical_policy_name(policy.py)做最终别名解析与合法性校验。合法的规范策略名定义于PolicyName枚举(policy.py),包括hipaa_safe_harbor、hipaa_expert_review_assist、gdpr_pseudonymization、gdpr_art9_health、research_limited_dataset、strict_no_leak、clinical_minimal_redaction、clinical_preserve、canada_pipeda、uk_ico_anonymisation、australia_privacy_act、china_pipl、india_dpdp_act、africa_malabo_baseline、za_popia、ng_ndpa、ke_dpa、india_health_id、eg_pdpl、ma_law_09_08等。
底层推理委托给openmed.processing.process_batch(batch.py),UDF 层固定注入的关键参数为:operation="deidentify"、policy=<规范化 profile>、method=<config.method>、confidence_threshold=<config.confidence_threshold>、loader=<worker 懒加载的 loader>、continue_on_error=False。这里continue_on_error=False与 postgres_plpython.py 中的continue_on_error=True形成对照:分布式 SQL 场景下任何一行失败都立即报错,避免把部分脱敏结果静默写回数据表。
懒加载与空输入:测试如何钉死这些行为
单元测试 test_distributed_sql_udf.py 用合成process_batch_fn与可计数的loader_factory验证了全部关键语义:
- 窗口化分批:
batch_size=2处理 3 行时,底层process_batch被调用两次,分别收到[2, 1]行;loader 工厂仅被调用 1 次——模型跨窗口复用得到直接验证。 - 空输入零开销:
udf([None, ""], ...)与udf.deidentify(None, ...)、udf.deidentify("", ...)均不触发 loader 工厂与process_batch(测试中 loader 工厂被设计为"被调用即抛AssertionError"),证明 NULL/空串透传不产生模型初始化与推理成本。 - profile 别名与分组:输入
["hipaa", "gdpr"]与"safe_harbor"时,底层按hipaa_safe_harbor、gdpr_pseudonymization、hipaa_safe_harbor三次调用分组推理,且全程只创建 1 个 loader——跨调用、跨 profile 复用得到验证。 - 行数对齐校验:
profiles长度与texts不等时抛ValueError("1 values for 2 text rows")。 - 描述符完整性:测试对
OPENMED_DEIDENTIFY_DESCRIPTOR做全量字典相等断言,锁定注册契约不漂移。
部署与 PHI 安全边界
该接入层面向"数据不离开集群"的部署形态,文档与源码对运维侧提出了明确要求:
- 完全离线部署:若部署必须完全离线,需在每个 worker 节点上预置 OpenMed 模型构件(默认
OpenMed/OpenMed-PII-SuperClinical-Small-44M-v1),并在接受查询前完成安装。worker 的默认 loader 工厂基于openmed.core.ModelLoader构建,其配置为缓存优先,不强制发起网络请求。 - 日志不落 PHI:该适配器自身不记录原始输入文本——这与 postgres_plpython.py 中"模块刻意不包含任何日志调用"的设计一脉相承。但引擎侧的查询日志、失败捕获(failure capture)与溢写(spill)配置需要运维方单独审查,确保它们不会把原始文本持久化到日志或临时文件。
- 错误消息脱敏:
_run_process_batch中任何行失败抛出的RuntimeError只包含行索引,不包含源文本与底层异常详情——这是防止引擎把 worker 异常文本转发到驱动端日志(driver logs)而泄漏 PHI 的刻意设计。
与同族集成的横向对照
分布式 SQL UDF 接入层只是 OpenMed 分布式去标识化家族的一员,在真实流水线中可按数据形态选择:
- Spark DataFrame:
openmed.interop.spark.SparkRedactionTransform提供显式的分区本地mapPartitions契约,见 spark.md;Spark 序列化的是不可变配置而非模型,每个分区尝试都会重建 worker-local 去标识器。 - Ray Data:
openmed.integrations.ray_map_batches.map_batches_deidentify提供有状态的Dataset.map_batches阶段,每个 actor 构造期加载一个模型管线并跨 batch 复用,见 ray-map-batches.md。 - PostgreSQL 内置:PL/Python 桥把 warm 模型管线缓存在数据库会话的
GD字典中,同一会话内标量与批量调用都不重复加载模型,见 postgres_plpython.py 及配套迁移脚本 postgres_deidentify.sql。
这些接入点的共同底层是openmed.processing.process_batch与openmed.core.ModelLoader:一份模型加载逻辑,多种引擎桥接形态。分布式 SQL UDF 接入层的定位,正是把这一能力以引擎无关的注册描述符形式暴露给那些没有专用适配器的 SQL 引擎。
总结
OpenMed 的分布式 SQL UDF 接入层把"逻辑标量、执行向量化"的openmed_deidentify函数包装成引擎无关的 Python 入口:注册描述符负责声明参数映射与向量化要求,DistributedSQLUDFConfig负责固化模型、策略、批大小等运行参数,DistributedSQLDeidentifyUDF负责懒加载、窗口切分、NULL/空串透传、profile 分组与行序对齐,模块级deidentify_batch则让注册入口零配置即可复用进程内 worker。接入方只需把描述符映射到目标引擎的 Python UDF 体系、在每个 worker 上预置模型构件,并审查引擎侧日志/溢写配置,即可在数据不离开集群的前提下,获得与 OpenMed 单机去标识化一致的 PHI 保护语义。
【免费下载链接】openmedLocal-first healthcare AI: clinical NER & HIPAA PII de-identification that runs 100% on-device. 2,200+ medical models, 21 languages, Apple MLX + Python, no cloud, no patient data leaving your network. Apache-2.0项目地址: https://gitcode.com/GitHub_Trending/ope/openmed
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考