ZeroFutureTech · AI Learning · 约 25 分钟阅读

Palantir Foundry:
从一堆文件到 Ontology,
到底发生了什么?

拆解 Foundry 数据集成的完整链路——从原始文件落地,到结构化清洗,到 Ontology 建模。用航空运营场景贯穿始终,所有技术细节标注官方出处。

多数 Palantir 的分析文章从 Ontology 讲起,告诉你"它是组织的数字孪生",然后就直接跳到 AIP Agent。但中间那段——一堆散乱的文件是怎么变成可操作的 Ontology 对象的——很少有人完整拆解。

这条链路才是 Palantir 真正的基本功。如果你不理解 pipeline,就无法判断 Ontology 的能力边界和部署成本。

全链路:五个阶段

我们用一个航空公司运营系统的场景来走这条链路。假设你是某航空公司的数字化负责人,手上有三类数据:结构化的航班运营 CSV、半结构化的乘客预订 JSON、非结构化的维修报告 PDF。

CSV / Parquet(结构化) JSON / XML(半结构化) PDF / 图片(非结构化) STAGE 1 Data Connection 200+ connectors · Batch / Streaming / CDC · 数据原封不动拉入 STAGE 2 Raw Datasets 版本化文件容器 · "Git for Data" · Dataset / Media Set 两种存储 Schema 推断 自动识别列 & 类型 JSON/XML 解析 嵌套结构 → 表格 AIP Doc Intelligence OCR + VLM → Markdown STAGE 3 STAGE 4 Curated Datasets 清洗、拼接、schema 验证后的标准表格数据 STAGE 5 · Ontology:Object Types → Properties → Links → Actions
1
Data Connection — 数据原封不动拉进来

Foundry 的第一原则是 as-is ingestion——数据从源系统拉进来时不做任何预处理。

"Foundry is also distinguished by a philosophy that data should be ingested 'as-is' from its most raw source, with no external preprocessing."— Data Connection Overview, Palantir Docs

为什么?因为一旦所有变换都在 Foundry 内部完成,从 Ontology 的任何一个 property value,都能顺着 pipeline 一路追溯到原始文件的某一行。这就是 end-to-end data lineage。

连接方式:通过 200+ 种连接器对接外部系统(数据库、ERP、API、文件系统、IoT 流),通过 Batch Sync(全量 SNAPSHOT 或增量 APPEND)将数据拉入 Foundry。

2
Raw Dataset — Foundry 的数据存储单元

在深入 pipeline 之前,有必要先理解 Foundry 里"数据"到底是怎么存的——这直接决定了后续 pipeline 能对它做什么操作。

Dataset:Foundry 的核心存储抽象

在 Foundry 里,Dataset 不是一张表,而是一个版本化的文件容器。它可以承载三类数据:

"Datasets are used to store and represent all types of data—structured, unstructured, and semi-structured: Structured (tabular) data consists of files in an open-source format such as Parquet, along with metadata about the columns. Unstructured datasets consist of files such as images, video, or PDFs, but do not have an associated schema. Semi-structured datasets contain file formats such as XML or JSON."— Core Concepts: Datasets, Palantir Docs
形态存储格式Schema航空场景
结构化Parquet 文件 + schema 元数据有(列名 + 类型)raw_flight_operations
半结构化JSON / XML 原始文件可推断,建议下游处理raw_passenger_bookings
非结构化原始二进制文件无(unstructured dataset)可存,但缺乏媒体处理工作流支持

Dataset 的每次变更都是一个原子性 Transaction,类比 Git 的 commit。四种事务类型:SNAPSHOT(全量替换当前视图)、APPEND(增量追加)、UPDATE(更新)、DELETE(删除)。这就是 Palantir 所说的 "Git for Data"——任何中间步骤出问题,都可以回溯到历史版本。

Media Set:面向媒体处理工作流的专用抽象

Dataset 技术上可以存非结构化文件,但在面向大规模文档、图片、音视频处理的场景下,工程实践中通常使用 Media Set——因为它在 Dataset 的存储基础之上,专门提供了 media item 管理、media reference 关联、媒体格式转换、以及后续 AI 提取工作流(Document Intelligence)的集成支持。

"A media set is a collection of media files with a common schema, for example, files of the same format. Media sets are designed to work with high-scale, unstructured data and enable the processing of media items such as audio, imagery, video, and documents."— Core Concepts: Media Sets, Palantir Docs

Media Set 有两个关键属性:

Schema Type(文件类别)——定义这个 Media Set 存什么类型的文件:document、image、audio、video、spreadsheet、dicom、email。创建时必须选定,之后只接受匹配类型的文件上传。

Primary Format(主格式)——指定文件的标准格式(如 PDF、PNG、MP4)。你也可以设定 additional input formats,上传时会自动转换为主格式。

Media Set 里的每个文件叫 Media Item,有唯一的 media_item_rid。关键设计是 Media Reference——你可以在 tabular dataset 里用一个特殊列引用 media item,把非结构化内容和结构化元数据关联起来:

"You can use media references to reference media set items in datasets. This is useful for associating media items with metadata or other information in a tabular format. For example, you can associate the original PDF with its file name, page count, and extracted text as additional columns."— Media Sets Overview, Palantir Docs

航空场景:维修报告 PDF 存储在一个 schema type 为 document、primary format 为 PDF 的 Media Set 里。后续 pipeline 提取出的结构化信息(报告编号、维修日期、零件列表等)存为普通 Dataset,通过 media_reference 列指向原始 PDF。

准确理解两者关系:Dataset 是 Foundry 的通用存储单元,可以承载结构化、半结构化、非结构化三类数据。Media Set 不是 Dataset 的替代,而是针对媒体处理工作流的专用上层抽象——当你需要处理大规模图片、视频、PDF 并接入 AI 提取管道时,Media Set 提供了 Dataset 本身没有的工作流支持。两者可以通过 Media Reference 关联。
3
Pipeline Transforms — 真正的"脏活"

目标:将所有数据统一变换为有 schema 的 tabular dataset。Foundry 提供两个主要的 pipeline 构建工具:

Pipeline Builder——低代码/无代码的可视化界面,拖拽式构建数据变换。两种操作类型:expressions(输入列→输出单列,如字符串拆分、类型转换)和 transforms(输入表→输出表,如 Filter、Pivot、Aggregate、Join)。

"Pipeline Builder uses a general model for describing data transformations. This backend is an intermediate layer between the tools used to write transformations and the execution of said transformations."— Pipeline Builder Overview, Palantir Docs

Code Repositories——用 Python(PySpark)或 Java 写 transform 代码。用 @transform_df 装饰器声明输入输出,返回 DataFrame 自动写入 output dataset。

根据数据类型,pipeline 走不同的路径:

路径 A:结构化数据 → Schema 自动推断 + 清洗

这是最简单的路径。CSV 或 JSON 文件已经是表格格式,Foundry 可以直接推断 schema:

"Foundry allows you to manually add a schema to datasets containing CSV or JSON files by selecting the Apply a schema button in the dataset. The Apply a schema button will automatically infer the schema based on a subset of the data."— Infer a schema, Palantir Docs

Schema 应用后,你可以进一步修改列类型、丢弃格式异常的行、添加额外列(如文件路径、导入时间戳、行号等)。

航空场景实例:航班运营数据清洗

假设 raw_flight_operations 是从 Oracle DB 同步过来的 CSV,包含 flight_id, tail_number, origin, destination, departure_time, status 等列。Pipeline 需要做的:

在 Pipeline Builder 中(无代码):选择 raw dataset 作为 Input → Apply schema(自动推断 departure_time 为 Timestamp 类型)→ Filter transform 过滤 status 不为 null 的行 → Cast expression 将 flight_id 统一为 String 类型 → 输出为 clean_flights。

在 Code Repositories 中(Python),同样的逻辑:

from transforms.api import transform_df, Input, Output
from pyspark.sql import functions as F

@transform_df(
    Output('/airlines/clean/flights'),
    raw=Input('/airlines/raw/flight_operations')
)
def clean_flights(raw):
    return (
        raw
        .filter(F.col('status').isNotNull())           # 过滤空状态
        .withColumn('flight_id',                       # 统一类型
            F.col('flight_id').cast('string'))
        .withColumn('departure_time',                  # 时间解析
            F.to_timestamp('departure_time',
                           'yyyy-MM-dd HH:mm:ss'))
        .dropDuplicates(['flight_id'])                # 去重
    )

@transform_df 装饰器做了三件事:声明输入 dataset 路径、声明输出 dataset 路径、把返回的 DataFrame 自动写入 output。这是 Foundry pipeline 的基本编程模型。

路径 B:半结构化数据 → JSON/XML 解析 + 摊平

半结构化数据的核心挑战是嵌套结构。一个 JSON 预订报文可能长这样:

{
  "booking_id": "BK-20250601-001",
  "passenger": {
    "name": "Zhang Wei",
    "passport": "E12345678"
  },
  "flights": [
    {"flight_id": "CA1234", "seat": "32A"},
    {"flight_id": "CA5678", "seat": "14C"}
  ]
}

Pipeline Builder 提供 parsing transform functions 把文件转为表格:

"Semi-structured data refers to a dataset that consists of files without a schema, making the data non-tabular. Pipeline Builder supports semi-structured data in the form of XML, JSON, and CSV files. You can use parsing transform functions to convert semi-structured files into tabular form and benefit from schema safety."— Pipeline Builder: Input datasets, Palantir Docs

航空场景实例:预订 JSON 摊平为两张表

一个 JSON 文件需要变成两个 output:

Output 1 — clean_passengers:提取 passenger 对象,取 booking_id 作为关联键。在 Pipeline Builder 中用 file transform 解析 JSON → 用 expression 提取嵌套字段 passenger.name、passenger.passport → 输出 tabular dataset。

Output 2 — clean_bookings:对 flights 数组使用 explode_array generator(Pipeline Builder 的内置功能),将数组展开为每行一条 flight 记录,附带 booking_id 和 passenger_id 作为外键。

在 Code Repositories 中,PySpark 的 explode 函数天然支持这种操作:

@transform_df(
    Output('/airlines/clean/bookings'),
    raw=Input('/airlines/raw/passenger_bookings')
)
def flatten_bookings(raw):
    # 读取 JSON → 展开 flights 数组 → 提取嵌套字段
    return (
        raw
        .select(
            F.col('booking_id'),
            F.col('passenger.name').alias('passenger_name'),
            F.col('passenger.passport').alias('passport'),
            F.explode('flights').alias('flight')      # 数组展开
        )
        .select(
            'booking_id', 'passenger_name', 'passport',
            F.col('flight.flight_id').alias('flight_id'),
            F.col('flight.seat').alias('seat')
        )
    )

一个 explode 调用就把嵌套的 JSON 数组"摊平"成了标准的 tabular 格式。这是处理半结构化数据最常用的模式。

路径 C:非结构化数据 → AIP Document Intelligence

这是最复杂的路径,也是 Palantir 近年投入最重的方向。非结构化数据(PDF、图片、音频)不能直接变成"行和列",需要内容提取。

AIP Document Intelligence 是 Palantir 在 AIP 层面专门为文档提取构建的产品级工具,提供两大类策略:

类别方法机制适用场景
TraditionalRaw text读取 PDF 元数据中的文本电子生成的 PDF
OCR传统 OCR 提取文本扫描件,版面不重要
Layout-aware OCROCR + bounding box,保留版面结构表格、多栏文档
Generative AIVLM 提取Vision Language Model 直接"看"页面,输出 Markdown复杂版面、混合内容
混合模式先 OCR 预处理,再将结果 + 页面图像传给 VLM最高精度场景
"Document preprocessing essentially runs traditional OCR on the document, then passes that output in addition to the document page itself to a VLM, giving the model more context to successfully analyze the document."— AIP Document Intelligence: Core Concepts, Palantir Docs

航空场景实例:维修报告 PDF 提取全流程

第一步——在 AIP Document Intelligence UI 中导入维修报告 Media Set,选择 Layout-aware OCR + VLM 混合策略,执行提取,查看 Markdown 输出(带 bounding box 对照原文)。

第二步——用 VLM-as-judge 自动评分,评估表格、标题、代码块等不同元素的提取质量。

第三步——满意后一键部署为 Python transform。系统自动生成代码仓库:

@transform.using(
    output=Output("ri.foundry.main.dataset.abc"),
    media_input=MediaSetInput("ri.mio.main.media-set.abc"),
    extractor=VisionLLMDocumentsExtractorInput(
        "ri.language-model-service..language-model.anthropic-claude-sonnet"
    ),
)
def extract(media_input, output, extractor):
    extracted_data = extractor.create_extraction(
        media_input,
        with_ocr=True,            # 混合模式:先 OCR 再 VLM
        prompt=USER_PROMPT,        # 可自定义提取 prompt
        thread_number=20           # 并行线程数
    )
    output.write_dataframe(extracted_data)
# 输出:media_item_rid | page_num | extraction_result (Markdown)

第四步——Document Intelligence 的输出是每页一行的 Markdown 文本。要从 Markdown 进一步提取结构化字段(aircraft_tail_number、maintenance_date、parts_replaced),还需要一个下游 pipeline——可以用 Pipeline Builder 的 LLM transform(分类、实体抽取),也可以写正则/NLP 代码。

第五步(容易被忽略)——LLM/VLM 提取的字段进入 curated dataset 之前,生产环境通常需要一套校验与人工复核闭环:

环节目的做法
Schema Validation确保提取字段类型和格式合规对 tail_number、date 格式做规则校验,不合规的行进异常队列
置信度过滤过滤低可信度的提取结果VLM 输出置信度低于阈值的记录标记为待复核
抽样人工复核评估整体提取质量随机抽取 N% 由人工对照原始 PDF 核验
异常队列处理边界情况和失败记录提取失败、格式异常的记录进人工处理队列

Palantir 的 Document Intelligence 官方也强调:用户在 UI 里先测试提取策略、查看 Markdown 与 bounding box 的对照、评估质量/速度/token cost,再部署。这个测试→评估→部署的闭环,本质上就是在做 offline 的 ground truth 校验,只是工具帮你做了一部分。

非结构化数据的完整路径(不是两跳,是五跳):
Media Set → Document Intelligence(OCR/VLM)→ Markdown/chunks → 字段提取(LLM/NLP)→ 校验 & 人工复核 → Curated Dataset → Ontology

Document Intelligence 解决了最难的第一跳(PDF → Markdown),但后续的提取精度、字段格式化、异常处理,每一步都有工程成本。这才是非结构化数据进入 Ontology 的真实代价。

对于音频数据,Foundry 的 Media Set transforms 内置了 Whisper 语音转录功能,可以直接用 transcribe 方法将音频转为文本。对于图片,支持 OCR、图像嵌入(image embeddings)、元数据提取(尺寸、分辨率等)。

从表格到 Ontology:四个核心概念

经过 pipeline 清洗,你有了干净的 tabular datasets。现在要把它们映射到 Ontology——Palantir 所谓的"组织的数字孪生"。

"The Foundry Ontology creates a complete picture of an organization's world by mapping datasets and models to object types, properties, link types, and action types."— Ontology Core Concepts, Palantir Docs

映射关系用 mental model 来理解很直接——dataset ≈ object type,row ≈ object,column ≈ property。这个类比适合入门,但实际建模时需要显式配置 backing datasource、property mapping、primary key / title key、link cardinality、property visibility,以及更复杂的情况如 Multi-datasource Object Types(MDO,一个 Object Type 从多个 datasource 取属性)。

1. Object Types — 每张表变成一类实体

创建 Object Type 时,你选择一个 backing dataset,系统自动把每一列映射为一个 property,你指定 primary key 和 title。

Backing DatasetObject TypePrimary KeyTitle 示例
clean_flightsFlightflight_id"CA1234"
clean_passengersPassengerpassenger_id"Zhang Wei"
clean_bookingsBookingbooking_idbooking_id
clean_aircraftAircrafttail_number"B-6075"
clean_maintenanceMaintenance Recordreport_idreport_id

2. Properties — 列变成属性

以 Flight 为例:flight_id → Flight ID (String, PK),tail_number → Tail Number (String, FK),departure_time → Departure Time (Timestamp),status → Status (String)。

支持的属性类型包括基础类型和高级类型:Vector(语义搜索)、Geopoint/Geoshape(地理)、Media Reference(关联原始 PDF/图片)、Time Series(传感器数据)、Struct(复杂嵌套结构)。

3. Link Types — 外键变成关系

Link 基于 foreign key ↔ primary key 匹配:

Aircraft Flight Booking Passenger Maintenance tail_number flight_id passenger_id tail_number(via Aircraft)

4. Action Types — 让 Ontology 可写

Actions 允许用户通过定义好的操作修改数据——创建对象、修改属性、建立/删除链接。

精妙设计:用户的编辑不直接写回 backing dataset,而是写入独立的 writeback dataset。原始数据永远不会被用户操作污染——你始终可以区分"系统同步的数据"和"人工修改的数据"。
Action Type操作涉及 Object业务逻辑
De-escalate AlertModifyFlight将 status 从 Alert 改为 Monitoring
Assign AircraftModify + LinkFlight, Aircraft修改 tail_number,自动建立链接
Create MaintenanceCreate + LinkMaintenance, Aircraft创建记录,关联 Aircraft
Rebook PassengerModify + RelinkBooking修改 flight_id,重建链接

Ontology 的逻辑层:Logic 和 Actions 是怎么来的

Ontology 映射完成后,你有了 Object Types、Properties、Link Types。但还缺两层关键能力:逻辑计算(基于数据自动推导新信息)和写操作(让业务用户可以修改数据)。这两层在 Palantir 里分别叫 Functions 和 Action Types。

Functions(逻辑计算层) TypeScript / Python 代码,运行在 Ontology 数据之上 Derived properties 派生属性 Aggregations 自定义聚合 Complex edits 批量修改逻辑 External queries 外部系统查询 Functions 可以被 Action Types 调用(function-backed actions) Action Types(写操作层) 定义用户可以对 Ontology 做的操作 Parameters(参数) Action 的输入:新角色、目标对象、日期… Rules(规则) Create / Modify / Delete 对象 Create / Delete Link Submission criteria(提交条件) 谁能执行?满足什么条件? Side effects(副作用) Notification 通知 Webhook 外部系统回调 触发其他 Pipeline Function-backed actions 复杂逻辑用代码实现 → Writeback Dataset(独立存储)

Functions:逻辑计算层

Functions 是运行在 Ontology 数据之上的 TypeScript / Python 代码,能做四件事:

能力说明航空场景例子
Derived Properties从现有属性实时计算出新属性delay_minutes = actual_departure - scheduled_departure
Aggregations跨对象的聚合计算avg_delay_per_airline = AVG(Flight.delay_minutes) GROUP BY airline
Complex Edits批量修改逻辑(配合 Action 使用)取消航班时自动释放所有座位、通知所有乘客
External Queries查询外部系统丰富 Ontology调用天气 API 给 Flight 对象附加当前天气数据
"Functions on objects (FOO) enables code authors to write business logic on object data and leverage this logic downstream in operational applications."— Functions on Objects: Overview, Palantir Docs

Action Types:写操作层(四个组件)

Action Type 不是简单的"让用户点编辑"——它是一套完整的写操作定义框架,由四个组件构成:

"An action type is the definition of a set of changes or edits to objects, property values, and links that a user can take at once. It also includes the side effect behaviors that occur with action submission."— Action Types: Overview, Palantir Docs
组件作用"改签乘客"实例
Parameters
参数
Action 的输入,类型化的表单字段passenger(Object Reference)、new_flight(Object Reference)、reason(enum)
Rules
规则
定义对 Ontology 的编辑操作Modify Booking.flight_id → new_flight;Delete Link → old_flight;Create Link → new_flight
Submission Criteria
提交条件
前置校验,不满足则无法提交new_flight.status ≠ Cancelled;new_flight.available_seats > 0;current_user.role ∈ [agent, supervisor]
Side Effects
副作用
执行后触发的外部效果Notification → passenger.email:"您的航班已改为 {new_flight}";Webhook → crew_management_api

Action 有两种实现方式:简单规则配置(UI 点选,不用写代码)和 function-backed action(用 Functions 层的代码实现复杂逻辑)。

"Rules define the logic of the action type that transform the parameters into Ontology edits. Create object: create an object of a predefined type. Modify object(s): modify an existing object. Delete link: delete a many-to-many link."— Action Types: Rules, Palantir Docs
Action 的执行时序:用户填入参数 → Submission Criteria 校验 → PASS:执行 Rules → 写入 Writeback Dataset → 触发 Side Effects → Ontology 立即更新 | FAIL:展示失败原因("目标航班已取消" / "无权限")。

数据怎么保持最新:增量更新机制

Ontology 建好之后,源系统的数据每天都在变化——新航班不断加入,维修记录持续更新,乘客随时改签。这些变化怎么反映到 Ontology 里?全量重算一遍代价太高,这就是增量更新要解决的问题。

增量更新贯穿整条链路三个层次:同步层 → Pipeline 层 → Ontology 索引层,每层都有独立的增量机制,必须端到端配合才能真正省时省力。

基础概念:APPEND vs SNAPSHOT

两种事务类型决定了整个增量机制能否成立:

APPEND 事务 只追加新文件,不修改已有文件 Week 1:100 万行(Transaction 1) Week 2:新增 10 万行(Transaction 2) Week 3:新增 8 万行(Transaction 3) ✓ 增量 pipeline 的基础 SNAPSHOT 事务 全量替换,覆盖当前视图所有数据 Week 1:100 万行(已覆盖) Week 2:110 万行(已覆盖) Week 3:118 万行(当前视图) ↻ 迫使下游全量重算
关键约束:UPDATE 事务(可能覆盖已有文件)会破坏 APPEND-only 要求,迫使下游所有 pipeline 回退到全量 SNAPSHOT 重算。只有在无法避免修改已有文件时才用 UPDATE。历史记录仍然保留(Git for Data 特性),但当前视图被替换。

三层增量更新:端到端配合

第 1 层:同步层 Data Connection — APPEND Sync 源系统有新数据 → 只拉增量行进 Foundry,不重传历史 ✓ 减轻源系统负担 第 2 层:Pipeline 层 @incremental() Transform read mode: added 只读新增行 write mode: modify 追加到已有输出 semantic_version bump 强制全量重算 第 3 层:Ontology 索引层 Funnel Service Changelog Job 自动计算 data diff Merge Job 合并 datasource + Actions Hydration 写入 OSv2 搜索节点 Ontology 可查询

增量不是一个开关:端到端契约

很容易产生一个误解:上游配了 APPEND sync,下游就自然是增量的。实际上,增量是一个端到端契约,链路上任何一个环节断裂,后面的增量就失效。

层次增量机制容易断裂的情况
Ingestion 层
Data Connection
JDBC incremental sync:基于严格递增的 cursor 字段(如 updated_at、auto-increment ID)只拉新行 源数据会更新已有行(非 insert-only);cursor 字段不严格单调递增;源系统无法提供可靠 change marker
Transform 层
@incremental() decorator
只读自上次 build 以来的新增事务(read mode: added) 上游 dataset 收到 UPDATE 事务(修改已有文件),下游 pipeline 被迫回退到 SNAPSHOT 全量重跑
Ontology 索引层
Funnel Service
Changelog job 计算 diff,只把变化量推入 OSv2 backing dataset 触发 SNAPSHOT,Funnel 需要全量重索引
两个常见陷阱:

1. 源数据可变(mutable)+ APPEND 事务 = 重复数据问题。如果源系统的记录会被修改(比如航班状态从 Scheduled 改为 Delayed),而你用 APPEND 事务同步,dataset 里就会同时存在同一条记录的旧版本和新版本。下游 transform 必须显式做 dropDuplicates 或取最新版本,否则 Ontology 里会出现脏数据。

2. File-based incremental UPDATE = 增量链断裂。File-based 的 UPDATE 事务(修改已有文件)虽然每次只摄入变更的文件,但会导致下游所有 pipeline 无法增量运行——因为输入 dataset 不再是 append-only,下游必须全量重跑来保证正确性。

设计增量链路的起点是:源系统能否提供可靠的 cursor 或 change marker(insert-only 或带 updated_at 的 upsert)?如果不能,从 ingestion 层就只能 SNAPSHOT,后续所有增量都是补救措施而非真正的端到端增量。

具体例子:航班数据增量更新

假设 clean_flights 有 2000 万行历史数据,每周新增 100 万行。如果每次全量重跑,每次要处理 2000 万行;用增量,每次只处理 100 万行——节省 95% 的计算量。

Pipeline Builder 中勾选 incremental 模式即可;Code Repositories 中用 @incremental() 装饰器:

from transforms.api import transform_df, Input, Output, incremental
from pyspark.sql import functions as F

@incremental(semantic_version=1, v2_semantics=True)
@transform_df(
    Output('/airlines/clean/flights'),
    raw=Input('/airlines/raw/flight_operations')
)
def clean_flights(ctx, raw):
    # 增量模式:raw 只包含自上次 build 以来的新增行
    # 无需改写任何逻辑,加装饰器即可生效
    return (
        raw
        .filter(F.col('status').isNotNull())
        .withColumn('departure_time',
            F.to_timestamp('departure_time', 'yyyy-MM-dd HH:mm:ss'))
        .dropDuplicates(['flight_id'])
    )
# semantic_version:改为 2 会强制全量重跑(逻辑变更时用)
@incremental() 参数作用何时用
semantic_version整数,bump 强制全量重算业务逻辑变了,历史数据需要重跑
snapshot_inputs指定某个 input 允许 UPDATE/DELETE部分输入不是 append-only
allow_retention允许 Foundry Retention 删除旧文件不需要永久保留历史数据
strict_append强制底层写入为 APPEND 事务需要严格保证 append-only 语义
v2_semantics=True启用新版增量语义所有新项目推荐默认开启

长期运行的隐患:小文件积累与 Projection Compaction

增量 pipeline 运行几周后,dataset 里会积累大量小 Parquet 文件——每次 APPEND 就是一批新文件。文件数达到数万、数十万时,读性能会明显下降。

积累 6 个月后(无 Projection) 数万个 2-5MB 小文件 → 读取慢,查询耗时长 → Spark 需要打开大量文件 → Projection Compaction 启用 Projection(自动 Compaction) 大文件 1(已排序 + hash-bucketed) 大文件 2 → 读性能与文件数无关 · 自动维护
Projection 工作方式:独立于 canonical dataset 的查询机制,自动把大量小文件合并成排序好的大文件(compaction)。Foundry 所有产品都知道怎么把 projection 和最新的 canonical 数据合并,读取结果和没有 projection 时完全一致——只是快很多。可以按 filter(优化过滤查询)或 join(优化关联查询,hash-bucketed)两种模式配置。

Ontology 层的增量:Link Type 和 Action Type

很多人只理解到 Object Type 的数据更新,忽略了 Link 和 Action 的同步机制——它们的更新方式完全不同:

Ontology 元素更新触发方式更新时机处理层
Object Type 属性 backing dataset 收到新事务 下次 Funnel batch pipeline 运行后 Funnel changelog + merge + hydration
One-to-many Link
(FK 型,如 Flight → Aircraft)
FK 属性值变化(随 Object 更新自动解析) 随 Object 更新,无需单独 sync 查询时动态解析,不需要 Funnel 单独处理
Many-to-many Link
(join table 型)
join table dataset 收到新事务 下次 Funnel batch pipeline 运行后 Funnel 独立为 link type 维护 changelog
Action Type 执行结果 用户触发 Action 立即写入 OSv2 live index Actions service → Funnel queue → 异步持久化

Action 的两段式更新是最关键的设计:

用户点击 "Rebook 改签" Actions Service 发送修改指令 立即 写入 OSv2 live index 异步 · 周期性 持久化到 Writeback Dataset Ontology 查询立即可见

这意味着:用户执行"Rebook 改签"后,Booking 对象的 flight_id 属性和 Link 链接立刻更新,不需要等 Funnel batch pipeline 的下次运行。但这次 Action 的持久化记录(写入 Writeback Dataset)是异步完成的——Funnel 有 offset 追踪保证故障恢复后能补齐未刷入的编辑。

现实问题:一张超大宽表怎么办?

前面的场景假设数据已范式化。但企业现实是:你拿到的往往是一张几百列的大宽表——ERP 导出的报表视图,order、customer、product、shipping 全打平在一张表里。

Palantir 的答案很明确:必须在 pipeline 层先拆分,没有捷径。

硬性约束:一个 Object Type 只能对应一个 backing dataset,一个 dataset 只能 back 一个 Object Type。

"Kitchen Sink" 反模式:Palantir 在 Best Practices 文档中专门警告——把所有列都映射为 properties 会让数据模型"被技术产物搞得一团糟"。只包含有明确业务含义且对 workflow 有用的列。

实际工程路径是三步:Pipeline 层拆分宽表为多个 focused datasets → 每个映射到独立 Object Type → 用 Link Types 关联。

关键瓶颈:谁来决定"怎么拆"?

Pipeline Builder 的 AIP 功能可以帮你写拆分代码。但决定怎么拆——200 列宽表拆成哪几个实体、order_status 属于 Order 还是应该独立——完全依赖人工判断,也就是 Palantir 的 Forward Deployed Engineer(FDE)。AIP Assist 是加速执行的 copilot,不是替代建模决策的 architect。

延伸思考:垂直行业的 Domain Ontology Template

上一节指出了 Palantir 模式的一个结构性成本:每个客户的 Ontology 设计都是从零开始的 FDE 定制项目。但这个问题在某些行业里,可能存在产品化的解法。

Palantir 有模板分发机制,但没做垂直行业模板

Palantir 的 Marketplace 支持把 object types 和 link types 打包成可复用产品,技术上完全可以做成一键安装的行业骨架。但作为平台厂商,Palantir 自己并没有针对特定行业做这件事。

机会在哪?

在业务实体高度标准化的行业——比如供应链(Shipment、Container、SKU、Warehouse、Carrier)、医疗(Patient、Encounter、Diagnosis、Medication、Provider)、金融(Account、Transaction、Instrument、Counterparty)——核心 Object Types 的跨企业共识度很高。这不是护城河。

真正值钱的不是"模板里有哪些 entity",而是:当客户给你一张源系统导出的 200 列宽表时,哪些列映射到哪个 entity。这套 mapping 规则才是 FDE 在每个客户那里反复花时间做的事。

产品化的方向不是一个静态模板,而是一套 source-system-aware 的 mapping 规则库。以供应链为例:

层级
内容
复用率
Layer 1
Domain Object Types
Shipment、Container、SKU、Warehouse 等标准 entity + Link Types
90%+
Layer 2
Source System Mappings
SAP MM/SD/WM、Oracle EBS、用友、金蝶等主流 ERP 标准表 → Object Type 映射
70%+
Layer 3
Customer Extensions
非标字段、定制业务逻辑、客户特有 entity
需定制

同样的三层架构可以套用到医疗(HL7 FHIR 标准 → 各 HIS 系统映射)、金融(SWIFT 消息标准 → 各核心系统映射)等领域。客户选了源系统类型,Layer 1 和 Layer 2 就大致有了,FDE 只需处理 Layer 3。

Palantir 没做这一步,因为它是平台厂商——没有动力把某个行业的 FDE 知识固化成产品。但对垂直行业玩家来说,这正好是杠杆点:谁能把这个领域知识标准化,谁就能在 Palantir 生态(或任何类似的 Ontology-driven 平台)里占据一个有结构性优势的位置。

底层真相:Ontology 到底"新"在哪?

如果你有数据库背景,看完可能会问:这不就是关系型 schema 加了一层业务语义吗?坦率说,从纯数据建模角度,确实如此。但 Ontology 真正的价值在三点:

第一,全链路可追溯。从 Ontology 任意 property value 到原始文件的某一行,整条 pipeline 有完整版本历史。合规行业的硬性需求。

第二,可写回。传统语义层是只读的。Ontology 通过 Action Types + Writeback Dataset 设计,让业务用户直接修改数据——且修改和原始数据分离存储。

第三,应用层直接消费。Workshop、Slate、OSDK、AIP 直接操作 Ontology 对象。业务用户点"Resolve Alert"按钮,背后自动完成 property 修改、link 更新、writeback 写入。

← AI 战略认知 · 全部文章