modules.mergers 模块帮助
本章节包含 modules.mergers 包中常用「数据合并/融合」模块的使用说明和示例,例如:
多表按索引字段做类似 SQL Join 的合并(
MergeTables)向已有表集合追加表并合并关系元数据(
AddTables)将输入表与 GDIM 数据库表对齐并生成可写入结构(
MergeGdimTables)多个
ResultModel/SingleResult混合合并为一个ResultModel``(``MergeResultModels)
MergeSingleResult
模块简介与适用场景
MergeSingleResult用于将多个SingleResult合并为一个SingleResult``(端口:``OutputSingleResult)。已不再推荐用于新流程——建议优先使用
MergeResultModels,因为SingleResult后续会逐步淘汰。典型适用场景:
多个模块分别产出不同 key 的结果,需要汇总成一个结果对象;
报告生成前,把多个计算模块的
SingleResult统一合并,便于后续写入模板/汇总输出;对冲突 key 做可控处理(合并策略见后文「合并规则与冲突策略」与「参数说明」)。
端口说明
输入端口 - 动态输入端口:类型为
SingleResult。需要先通过add_dynamic_ports_in("InputXXX")显式添加端口后,才能在 pipeline 里连接到该端口。输出端口 -
OutputSingleResult:合并后的SingleResult;当未连接任何输入端口、或(在all_ports_required=True时)存在未就绪端口时为None。
如何添加动态输入端口(必读)
MergeSingleResult 的输入端口是 “动态端口”,也就是端口数量不固定。使用时你需要:
先在模块上用
add_dynamic_ports_in("InputXXX")创建若干输入端口;再把上游模块的
OutputSingleResult分别连接到这些端口。
Note
动态输入端口名必须以
Input开头(例如InputSR1/InputA),否则会报错。默认
all_ports_required=True:只要你创建过的任一动态端口数据为None,该模块就会输出None;若某些分支是可选的,请把all_ports_required=False。
合并规则与冲突策略(核心概念)
对于不同 key(不冲突):像合并字典一样直接合并。
对于同 key(冲突):按策略合并冲突值;全局默认策略由
default_strategy控制,特定 key 可通过key_strategies单独指定(详见「参数说明」)。内置策略:
overwrite_first、overwrite_last``(默认)、``concat_list、concat_string。
快速上手示例:默认策略合并(最后一个覆盖)
from gdisdk.modules.mergers import MergeSingleResult
m = MergeSingleResult(mname="MergeSR")
m.add_dynamic_ports_in("InputSR1")
m.add_dynamic_ports_in("InputSR2")
# m.InputSR1 = sr1 # SingleResult
# m.InputSR2 = sr2
m.execute()
merged = m.OutputSingleResult.data
快速上手示例:为不同 key 指定不同策略
from gdisdk.modules.mergers import MergeSingleResult
m = MergeSingleResult(mname="MergeSR")
m.default_strategy = "overwrite_last"
m.key_strategies = {
"project_name": "concat_string", # 也可以用 key 的 title 来匹配
"bore_ids": "concat_list",
}
m.string_separator = " | "
m.add_dynamic_ports_in("InputSR1")
m.add_dynamic_ports_in("InputSR2")
# m.InputSR1 = sr1
# m.InputSR2 = sr2
m.execute()
merged = m.OutputSingleResult.data
参数说明
参数名 |
类型 |
默认值 |
说明 |
|---|---|---|---|
|
|
|
为 |
|
|
|
冲突 key 的默认合并策略。 |
|
|
|
|
|
|
|
为特定 key 指定策略的字典;key 可以写结果的 |
Note
合并冲突 key 时,UnitResult 的元信息(如 title/unit/description)默认沿用“第一个出现的结果”的元信息。
在 pipeline 中的使用方式
from gdisdk.pipeline.pipeline import PipeLine
from gdisdk.modules.mergers import MergeSingleResult
pipe = PipeLine(app_name="MergeSingleResultDemo", app_title="合并 SingleResult 示例")
m = MergeSingleResult(mname="MergeSR")
m.all_ports_required = False
m.add_dynamic_ports_in("InputSR1")
m.add_dynamic_ports_in("InputSR2")
m.add_dynamic_ports_in("InputSR3")
# links = (
# mod1.OutputSingleResult >> m.InputSR1
# | mod2.OutputSingleResult >> m.InputSR2
# | mod3.OutputSingleResult >> m.InputSR3
# )
# pipe.add_links(links)
# pipe.add_module(m)
# pipe.run()
更多信息
MergeResultModels
模块简介与适用场景
MergeResultModels用于将多个ResultModel和/或SingleResult合并为一个ResultModel``(端口:``OutputResultModel)。推荐在新流程中优先使用本模块,尤其适合从
SingleResult逐步迁移到ResultModel的场景。典型适用场景:
多个模块输出
ResultModel,需要汇总成一个统一结果对象供下游消费;老流程中仍有
SingleResult输出,但下游希望统一接收ResultModel;需要对同名字段做冲突处理(合并策略见后文「字段展开与合并规则」与「参数说明」)。
端口说明
输入端口 - 动态输入端口:类型支持
ResultModel和SingleResult。需要先通过add_dynamic_ports_in("InputXXX")显式添加端口后,才能在 pipeline 里连接到该端口。输出端口 -
OutputResultModel:合并后的ResultModel;当未连接任何输入端口、或(在all_ports_required=True时)存在未就绪端口时为None。
如何添加动态输入端口(必读)
MergeResultModels 的输入端口是“动态端口”,也就是端口数量不固定。使用时你需要:
先在模块上用
add_dynamic_ports_in("InputXXX")创建若干输入端口;再把上游模块的
OutputResultModel或OutputSingleResult分别连接到这些端口。
Note
动态输入端口名必须以
Input开头(例如InputRM1/InputSR1/InputA),否则会报错。默认
all_ports_required=True:只要你创建过的任一动态端口数据为None,该模块就会输出None;若某些分支是可选的,请把all_ports_required=False。
字段展开与合并规则(核心概念)
对于
SingleResult输入:会按UnitResult.name展开为字段名,UnitResult.value作为字段值;若要在key_strategies中按标题匹配,可使用UnitResult.title。对于
ResultModel输入:会按model_dump()结果展开字段;若字段定义带有Field(title=...),则标题也可用于key_strategies匹配。对于不同字段名(不冲突):像合并字典一样直接合并。
对于同字段名(冲突):按策略合并冲突值;全局默认策略由
default_strategy控制,特定字段可通过key_strategies单独指定(详见「参数说明」)。内置策略:
overwrite_first、overwrite_last``(默认)、``concat_list、concat_string。
快速上手示例:合并多个 ResultModel
from pydantic import BaseModel, Field
from gdisdk.modules.mergers import MergeResultModels
class ProjectInfo(BaseModel):
project_name: str = Field(title="项目名称")
section_name: str = Field(title="标段名称")
class ProjectStats(BaseModel):
bore_count: int = Field(title="钻孔数量")
m = MergeResultModels(mname="MergeRM")
m.add_dynamic_ports_in("InputRM1")
m.add_dynamic_ports_in("InputRM2")
# m.InputRM1 = rm1 # ResultModel
# m.InputRM2 = rm2
m.execute()
merged = m.OutputResultModel.data
快速上手示例:混合合并 ResultModel 和 SingleResult
from gdisdk.modules.mergers import MergeResultModels
m = MergeResultModels(mname="MergeMixedResults")
m.all_ports_required = False
m.default_strategy = "overwrite_last"
m.key_strategies = {
"备注": "concat_string", # 可按 title 匹配
"bore_ids": "concat_list", # 也可按字段名匹配
}
m.string_separator = " | "
m.add_dynamic_ports_in("InputRM1")
m.add_dynamic_ports_in("InputSR1")
# m.InputRM1 = rm1
# m.InputSR1 = sr1
m.execute()
merged = m.OutputResultModel.data
参数说明
参数名 |
类型 |
默认值 |
说明 |
|---|---|---|---|
|
|
|
为 |
|
|
|
冲突字段的默认合并策略。 |
|
|
|
|
|
|
|
为特定字段指定策略的字典;key 可以写字段名,也可以写字段标题(对 |
Note
输出是运行时动态创建的 Pydantic
BaseModel实例,其字段集合是所有输入字段名的并集。若只有一个有效输入,模块仍会输出一个
ResultModel;不会回传原始SingleResult类型。当多个输入存在同名字段时,冲突判断基于“字段名”本身;
title只用于key_strategies的匹配。
在 pipeline 中的使用方式
from gdisdk.pipeline.pipeline import PipeLine
from gdisdk.modules.mergers import MergeResultModels
pipe = PipeLine(app_name="MergeResultModelsDemo", app_title="合并 ResultModel 示例")
m = MergeResultModels(mname="MergeResults")
m.all_ports_required = False
m.add_dynamic_ports_in("InputRM1")
m.add_dynamic_ports_in("InputRM2")
m.add_dynamic_ports_in("InputSR1")
# links = (
# mod_result_1.OutputResultModel >> m.InputRM1
# | mod_result_2.OutputResultModel >> m.InputRM2
# | mod_single.OutputSingleResult >> m.InputSR1
# )
# pipe.add_links(links)
# pipe.add_module(m)
# pipe.run()
# merged = m.OutputResultModel.data
更多信息
MergeTables
模块简介与适用场景
MergeTables将TableCollection内多张TableData按索引字段合并为一张表(类似 SQL Join);输入为单张TableData时原样输出。典型适用场景:多表横向拼接字段(如基本信息表 + 计算结果表 + 备注表)。
join_key、tables_to_merge等配置见「参数说明」。
端口说明
输入端口 -
InputTables:输入表集合(TableCollection)或单表(TableData)。输出端口 -
OutputTable:合并后的TableData;当输入为空、或无法找到可用 join 键/可合并表为空时为None。
快速上手示例:合并集合内全部表(自动选择 join_key)
from gdisdk.modules.mergers import MergeTables
m = MergeTables(mname="MergeTables")
# m.InputTables = tables # TableCollection
m.execute()
out_table = m.OutputTable.data
快速上手示例:指定 join_key(可写列名或字段标题)
from gdisdk.modules.mergers import MergeTables
m = MergeTables(mname="MergeByKey")
m.join_key = "bore_number" # 列名
# m.join_key = "钻孔编号" # 字段标题(若能映射到列名也可)
# m.InputTables = tables
m.execute()
out_table = m.OutputTable.data
参数说明
参数名 |
类型 |
默认值 |
说明 |
|---|---|---|---|
|
|
|
合并索引字段(列名或字段标题)。为 |
|
|
|
指定要参与合并的表(表名或表标题)。为 |
|
|
|
合并方式(pandas merge 语义)。常用: |
|
|
|
重名列后缀(按 pandas merge 语义)。 |
|
|
|
是否对结果按 join key 排序。 |
|
|
|
是否尽量避免复制(按 pandas merge 语义)。 |
|
|
|
是否添加来源指示列(按 pandas merge 语义)。 |
|
|
|
关系校验(按 pandas merge 语义);打开可用于发现重复/主键不唯一等数据质量问题,但会增加开销。 |
|
|
|
是否忽略索引并重新编号(按 pandas merge 语义)。 |
Note
输入为
TableData时,该表会被直接返回(等价于“未做合并”)。输入为
TableCollection时会按顺序逐张合并;若某张表不包含 join 键会被跳过并给出警告。
在 pipeline 中的使用方式
from gdisdk.pipeline.pipeline import PipeLine
from gdisdk.modules.mergers import MergeTables
pipe = PipeLine(app_name="MergeTablesDemo", app_title="多表合并示例")
m = MergeTables(mname="MergeTables")
m.how = "left"
m.join_key = "bore_number"
# links = upstream.OutputTables >> m.InputTables
# pipe.add_links(links)
# pipe.add_module(m)
# pipe.run()
# out = m.OutputTable.data
更多信息
MergeGdimTables
模块简介与适用场景
MergeGdimTables将输入表与 GDIM 表按匹配键对齐,输出结构以 GDIM 表 schema 为准,并填充gdim_id,供GdimTableWriter写入前预处理。匹配方式:主表用自身主键(1:1 更新/插入);子表用父表主键(1:N 整组替换);也可设
join_key。join_key在两侧都唯一时按主表逻辑,任一侧有重复时按子表逻辑。GDIM 中存在、输入中不存在的键**不会**出现在输出中,因而也不会被删除。
输出列结构始终跟 GDIM 表走,而不是跟输入表走。1:1
overwrite更新行会保留输入中没有的 GDIM 列原值;插入行和 1:N 替换行的这些列填None。
端口说明
输入端口 -
InputTables:待对齐合并的输入表(TableCollection或TableData)。 -InputGdimTables:作为基准的 GDIM 表集合(TableCollection),每张表必须包含gdim_id列。输出端口 -
OutputTables:合并后的表(TableCollection或TableData,与输入类型保持一致)。
输出行为
1:1(主表,或 ``join_key`` 两侧均唯一)
1:N(子表,或 ``join_key`` 任一侧不唯一)
ignore:该组键只要在 GDIM 中已存在,整组跳过;只插入全新组。overwrite:匹配组的旧 GDIM 行输出为墓碑行;全部输入行作为新行插入(gdim_id=None)。
墓碑行(1:N overwrite 的删除约定)
墓碑行与插入行在**同一张表**里:完整 GDIM schema,
gdim_id保留,其余列全部为None。GdimTableWriter/write_table_data会先把这类行识别为删除标记,调用GDIM删除行接口,再写入其余插入/更新行。
快速上手示例:只输出新增行(ignore,默认)
from gdisdk.modules.mergers import MergeGdimTables
m = MergeGdimTables(mname="MergeWithGdim")
m.how = "ignore"
# m.InputTables = input_tables # TableCollection 或 TableData
# m.InputGdimTables = gdim_tables # TableCollection(含 gdim_id)
m.execute()
out = m.OutputTables.data
快速上手示例:允许更新(overwrite)并启用自动列映射
from gdisdk.modules.mergers import MergeGdimTables
m = MergeGdimTables(mname="MergeWithGdimOverwrite")
m.how = "overwrite"
m.auto_map_columns = True
# m.InputTables = input_tables
# m.InputGdimTables = gdim_tables
m.execute()
out = m.OutputTables.data
快速上手示例:按钻孔整组替换子表(join_key 不唯一 → 1:N)
from gdisdk.modules.mergers import MergeGdimTables
m = MergeGdimTables(mname="ReplaceByBore")
m.how = "overwrite"
m.auto_map_columns = True
m.join_key = "bore_number" # 同一钻孔多行,按组替换而不是按 depth 更新
# m.InputTables = input_tables # 已映射到目标表结构的新数据
# m.InputGdimTables = gdim_tables # 现有 GDIM 子表(含 gdim_id)
m.execute()
out = m.OutputTables.data
# 输出 = 匹配孔的墓碑行 + 全部输入行(gdim_id=None)
参数说明
参数名 |
类型 |
默认值 |
说明 |
|---|---|---|---|
|
|
|
|
|
|
|
是否自动做列映射。为 |
|
|
|
匹配列名或标题。两侧都唯一时按主表 1:1;任一侧重复时按子表 1:N。未设置时主表用自身主键,子表用父表主键。 |
Note
在 pipeline 中的使用方式
from gdisdk.pipeline.pipeline import PipeLine
from gdisdk.modules.mergers import MergeGdimTables
pipe = PipeLine(app_name="MergeGdimTablesDemo", app_title="与 GDIM 表对齐示例")
m = MergeGdimTables(mname="MergeWithGdim")
m.how = "ignore"
# links = input_reader.OutputTables >> m.InputTables | gdim_reader.OutputTables >> m.InputGdimTables
# pipe.add_links(links)
# pipe.add_module(m)
# pipe.run()
# out = m.OutputTables.data
更多信息
AddTables
模块简介与适用场景
AddTables用于把待追加的表合并进已有TableCollection,或从单张TableData出发构建集合并继续追加(端口:OutputTables)。与
MergeTables不同:本模块在**集合层面**追加整张表,不做按 join 键的行级合并。关系参数(
main_table_name/sub_table_names/primary_key)会**合并**进已有TableCollection结构,而不是整体覆盖。典型适用场景:
多个读取模块分别产出表,需要汇总到一个
TableCollection供下游统一消费;在已有集合上追加新表,同时补充主表、子表或主键字段等关系元数据;
从单张
TableData起步,逐步组装多表集合后再交给MergeTables/MergeGdimTables等模块。
端口说明
输入端口 -
InputTablesToAdd:待追加的表,类型为TableData或TableCollection。 -InputTablesAddedTo:目标集合或基准表,类型为TableCollection或TableData;若为单张TableData,会先包装为新的TableCollection。输出端口 -
OutputTables:追加后的TableCollection;当InputTablesToAdd或InputTablesAddedTo任一未就绪时为None。
合并规则(核心概念)
表数据:在
InputTablesAddedTo的副本上追加InputTablesToAdd中的全部表(TableCollection会逐表追加,TableData追加单表)。主表(``main_table``):
main_table_name可写表名或表标题;若与已有主表不同,会合并为list[str],不会替换已有主表。子表(``sub_tables``):
sub_table_names可写表名或表标题。若为list,条目会合并进现有sub_tables``(``list或dict[str, list[str]]);若为dict,键为主表名/标题,值为子表名/标题列表(见TableCollection.sub_tables)。主键(``primary_key``):
primary_key可写字段名或字段标题,按主表解析后合并进现有primary_key``(``str或dict[str, str])。集合元数据:
collection_metadata控制输出TableCollection的TableMetadata来源;详见参数说明。
快速上手示例:向已有集合追加表
from gdisdk.modules.mergers import AddTables
m = AddTables(mname="AddTables")
m.InputTablesAddedTo = existing_collection # TableCollection
m.InputTablesToAdd = new_table # TableData 或 TableCollection
m.execute()
out = m.OutputTables.data
快速上手示例:追加表并更新关系元数据
from gdisdk.modules.mergers import AddTables
m = AddTables(mname="AddTablesWithMeta")
m.InputTablesAddedTo = base_collection
m.InputTablesToAdd = extra_tables
m.main_table_name = "钻孔基本信息表" # 表名或表标题
m.sub_table_names = ["试验成果表", "地层描述表"]
m.primary_key = "钻孔编号" # 字段名或字段标题
m.collection_metadata = "combined"
m.combine_sep = "\n"
m.execute()
out = m.OutputTables.data
参数说明
参数名 |
类型 |
默认值 |
说明 |
|---|---|---|---|
|
|
|
追加为主表的表名或表标题;与已有 |
|
|
|
子表名或表标题; |
|
|
|
主键字段名或字段标题;按主表解析后合并进现有 |
|
|
|
输出集合元数据来源:直接传入 |
|
|
|
|
Note
若
InputTablesToAdd或InputTablesAddedTo未连接/未赋值,OutputTables为None。表名、表标题、字段名、字段标题在关系参数中均可使用;模块会按
TableCollection/TableData的实际结构解析。构造函数中的
tables_to_add/tables_added_to仅作端口初值;文档示例推荐通过InputTablesToAdd/InputTablesAddedTo赋值。
在 pipeline 中的使用方式
from gdisdk.pipeline.pipeline import PipeLine
from gdisdk.modules.mergers import AddTables
pipe = PipeLine(app_name="AddTablesDemo", app_title="追加表到集合示例")
m = AddTables("AddTables")
m.main_table_name = "钻孔基本信息表"
m.collection_metadata = "tables_added_to"
# links = (
# base_reader.OutputTables >> m.InputTablesAddedTo
# | extra_reader.OutputTables >> m.InputTablesToAdd
# )
# pipe.add_links(links)
# pipe.add_module(m)
# pipe.run()
# out = m.OutputTables.data
更多信息