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_firstoverwrite_last``(默认)、``concat_listconcat_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

参数说明

MergeSingleResult 参数一览

参数名

类型

默认值

说明

all_ports_required

bool

True

True 时,动态输入端口只要有任一端口数据为 None 就不会执行合并(输出 None);为 False 时会尽量合并已有的输入。

default_strategy

Literal["overwrite_first","overwrite_last","concat_list","concat_string"]

"overwrite_last"

冲突 key 的默认合并策略。overwrite_first / overwrite_last 取第一个或最后一个输入值;concat_list 拼成 list(若输入值本身为 list 会自动展开,避免嵌套);concat_stringstring_separator 拼成字符串(若值为 list 会先把元素转为字符串)。

string_separator

str

", "

concat_string 策略的拼接分隔符。

key_strategies

dict[str, Literal["overwrite_first","overwrite_last","concat_list","concat_string"]] | None

None

为特定 key 指定策略的字典;key 可以写结果的 nametitle;未指定的 key 使用 default_strategy

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

    • 需要对同名字段做冲突处理(合并策略见后文「字段展开与合并规则」与「参数说明」)。

端口说明

  • 输入端口 - 动态输入端口:类型支持 ResultModelSingleResult。需要先通过 add_dynamic_ports_in("InputXXX") 显式添加端口后,才能在 pipeline 里连接到该端口。

  • 输出端口 - OutputResultModel:合并后的 ResultModel;当未连接任何输入端口、或(在 all_ports_required=True 时)存在未就绪端口时为 None

如何添加动态输入端口(必读)

MergeResultModels 的输入端口是“动态端口”,也就是端口数量不固定。使用时你需要:

  • 先在模块上用 add_dynamic_ports_in("InputXXX") 创建若干输入端口;

  • 再把上游模块的 OutputResultModelOutputSingleResult 分别连接到这些端口。

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_firstoverwrite_last``(默认)、``concat_listconcat_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

参数说明

MergeResultModels 参数一览

参数名

类型

默认值

说明

all_ports_required

bool

True

True 时,动态输入端口只要有任一端口数据为 None 就不会执行合并(输出 None);为 False 时会尽量合并已有的输入。

default_strategy

Literal["overwrite_first","overwrite_last","concat_list","concat_string"]

"overwrite_last"

冲突字段的默认合并策略。overwrite_first / overwrite_last 取第一个或最后一个输入值;concat_list 拼成 list(若输入值本身为 list 会自动展开,避免嵌套);concat_stringstring_separator 拼成字符串(若值为 list 会先把元素转为字符串)。

string_separator

str

", "

concat_string 策略的拼接分隔符。

key_strategies

dict[str, Literal["overwrite_first","overwrite_last","concat_list","concat_string"]] | None

None

为特定字段指定策略的字典;key 可以写字段名,也可以写字段标题(对 SingleResultUnitResult.title,对 ResultModelField(title=...));未指定的字段使用 default_strategy

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

模块简介与适用场景

  • MergeTablesTableCollection 内多张 TableData 按索引字段合并为一张表(类似 SQL Join);输入为单张 TableData 时原样输出。

  • 典型适用场景:多表横向拼接字段(如基本信息表 + 计算结果表 + 备注表)。join_keytables_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

参数说明

MergeTables 参数一览

参数名

类型

默认值

说明

join_key

str | None

None

合并索引字段(列名或字段标题)。为 None 时,会自动选择“所有待合并表都包含的第一个公共列”作为 join 键;不包含 join 键的表会被跳过。

tables_to_merge

list[str] | str | None

None

指定要参与合并的表(表名或表标题)。为 None 时考虑集合内所有表;若为 str 会自动转为单元素列表。

how

Literal["left","right","outer","inner"]

"inner"

合并方式(pandas merge 语义)。常用:inner 仅保留匹配行;left 保留左表全部行。

suffixes

tuple[str, str]

("_x","_y")

重名列后缀(按 pandas merge 语义)。

sort

bool

False

是否对结果按 join key 排序。

copy

bool

True

是否尽量避免复制(按 pandas merge 语义)。

indicator

bool

False

是否添加来源指示列(按 pandas merge 语义)。

validate

Literal["one_to_one","one_to_many","many_to_one","many_to_many"] | None

None

关系校验(按 pandas merge 语义);打开可用于发现重复/主键不唯一等数据质量问题,但会增加开销。

ignore_index

bool

False

是否忽略索引并重新编号(按 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_keyjoin_key 在两侧都唯一时按主表逻辑,任一侧有重复时按子表逻辑。

  • GDIM 中存在、输入中不存在的键**不会**出现在输出中,因而也不会被删除。

  • 输出列结构始终跟 GDIM 表走,而不是跟输入表走。1:1 overwrite 更新行会保留输入中没有的 GDIM 列原值;插入行和 1:N 替换行的这些列填 None

端口说明

  • 输入端口 - InputTables:待对齐合并的输入表(TableCollectionTableData)。 - InputGdimTables:作为基准的 GDIM 表集合(TableCollection),每张表必须包含 gdim_id 列。

  • 输出端口 - OutputTables:合并后的表(TableCollectionTableData,与输入类型保持一致)。

输出行为

  • 1:1(主表,或 ``join_key`` 两侧均唯一)

    • ignore:剔除已存在键,只输出新键,gdim_id=None

    • overwrite:匹配行带上原 ``gdim_id``(更新),新键 ``gdim_id=None``(插入)。

  • 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)

参数说明

MergeGdimTables 参数一览

参数名

类型

默认值

说明

how

Literal["ignore","overwrite"]

"ignore"

ignore 不碰已有键;overwrite 在 1:1 时更新匹配行,在 1:N 时对匹配组输出墓碑行并插入全部输入行。

auto_map_columns

bool

False

是否自动做列映射。为 True 时会做更灵活的 name/title 交叉匹配(更慢);为 False 时只做快速匹配(列名完全一致或列标题完全一致)。

join_key

str | list[str] | None

None

匹配列名或标题。两侧都唯一时按主表 1:1;任一侧重复时按子表 1:N。未设置时主表用自身主键,子表用父表主键。

Note

  • 若既没有 join_key、集合里也解析不出主键,则所有 gdim_id 置为 ``None``(按新增处理)。

  • InputGdimTables 中的每张表都必须包含 gdim_id 列,否则会抛出 ValueError

  • 若某张输入表在 InputGdimTables 中找不到同名/同标题的 GDIM 表,则该表会原样返回(不会添加 gdim_id,也不会做列映射)。

  • InputGdimTables 为空集合,则直接返回 ``InputTables``(不做处理)。

  • 1:N overwrite 的墓碑行由 GdimTableWriter 删除;模块本身不调用删除接口。

在 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:待追加的表,类型为 TableDataTableCollection。 - InputTablesAddedTo:目标集合或基准表,类型为 TableCollectionTableData;若为单张 TableData,会先包装为新的 TableCollection

  • 输出端口 - OutputTables:追加后的 TableCollection;当 InputTablesToAddInputTablesAddedTo 任一未就绪时为 None

合并规则(核心概念)

  • 表数据:在 InputTablesAddedTo 的副本上追加 InputTablesToAdd 中的全部表(TableCollection 会逐表追加,TableData 追加单表)。

  • 主表(``main_table``)main_table_name 可写表名或表标题;若与已有主表不同,会合并为 list[str],不会替换已有主表。

  • 子表(``sub_tables``)sub_table_names 可写表名或表标题。若为 list,条目会合并进现有 sub_tables``(``listdict[str, list[str]]);若为 dict,键为主表名/标题,值为子表名/标题列表(见 TableCollection.sub_tables)。

  • 主键(``primary_key``)primary_key 可写字段名或字段标题,按主表解析后合并进现有 primary_key``(``strdict[str, str])。

  • 集合元数据collection_metadata 控制输出 TableCollectionTableMetadata 来源;详见参数说明。

快速上手示例:向已有集合追加表

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

参数说明

AddTables 参数一览

参数名

类型

默认值

说明

main_table_name

str | None

None

追加为主表的表名或表标题;与已有 main_table 合并,不替换。

sub_table_names

list[str] | dict[str, list[str]] | None

None

子表名或表标题;list 时合并进现有 sub_tablesdict 时按主表分组合并。

primary_key

str | None

None

主键字段名或字段标题;按主表解析后合并进现有 primary_key

collection_metadata

TableMetadata | dict | Literal["tables_to_add","tables_added_to","combined"]

"tables_added_to"

输出集合元数据来源:直接传入 TableMetadata/dict 时覆盖对应字段;"tables_to_add" / "tables_added_to" 取自对应输入;"combined" 时用 combine_sep 拼接两侧非空字段。

combine_sep

str

"\\n"

collection_metadata="combined" 时,合并元数据字段值的分隔符。

Note

  • InputTablesToAddInputTablesAddedTo 未连接/未赋值,OutputTablesNone

  • 表名、表标题、字段名、字段标题在关系参数中均可使用;模块会按 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

更多信息