Data Formulator 自定义数据源 Loader 插件开发指南:ExternalDataLoader 接口、插件注册与连接配置全解析

发布时间:2026/9/13 13:02:48
Data Formulator 自定义数据源 Loader 插件开发指南:ExternalDataLoader 接口、插件注册与连接配置全解析 Data Formulator 自定义数据源 Loader 插件开发指南ExternalDataLoader 接口、插件注册与连接配置全解析【免费下载链接】data-formulator Data Formulator is an interactive AI-powered data analysis system makes it easy to connect, explore and visualize data.项目地址: https://gitcode.com/GitHub_Trending/da/data-formulator本指南适用于 Data Formulator 0.7。面向需要接入新数据库、报表系统、对象存储或内部数据服务的开发者与管理员讲解如何通过ExternalDataLoader接口以零源码修改的方式扩展数据源以及如何通过DataConnector自动获得连接、浏览、预览、导入与刷新能力。读完本文你将能够编写一个可注册、可配置、可安全用于生产环境的自定义数据源 Loader 插件。导读Data Formulator 0.7 将数据源扩展统一收敛为ExternalDataLoader DataConnector模型开发者只需实现一个继承ExternalDataLoader的适配器子类框架会自动把它包装成可连接、可浏览、可预览、可导入、可刷新的连接实例connector instance并统一暴露/api/connectors/...API。本文以官方插件文档为主体结合仓库源码与示例插件完整覆盖扩展模型、两种接入场景、Loader 必须实现与推荐实现的方法、连接实例配置、认证与凭证、安全要求及验证清单帮助你正确接入 PostgreSQL、MySQL、Superset、S3、BigQuery 之外的全新数据源类型。1. 当前扩展模型ExternalDataLoader 与 DataConnectorData Formulator 0.7 使用统一的扩展模型ExternalDataLoader # 具体数据源实现的适配器接口 DataConnector # 框架自动提供的连接实例管理层 connectors.yaml # 管理员预配置连接实例 DF_PLUGIN_DIR # 不改源码外加 loader type 的目录各组件职责如下组件职责ExternalDataLoader定义某一种数据源如何连接、如何列举对象、如何读取数据的抽象基类external_data_loader.pyDataConnector自动包装任意 loader 类生成带认证、catalog、数据路由的 Flask Blueprintdata_connector.pyconnectors.yaml管理员预置的全局连接实例配置文件位于DATA_FORMULATOR_HOME/下DF_PLUGIN_DIR存放外部 loader 插件文件*_data_loader.py的目录新增数据源时通常只需要实现一个ExternalDataLoader子类。DataConnector会通过DataConnector.from_loader(...)自动将其包装为 connector instance并在启动时注册一组统一路由见 data_connector.pyGET/POST /api/connectors列举 / 创建连接实例POST /api/connectors/get-catalog、get-catalog-tree、search-catalog目录浏览与搜索POST /api/connectors/preview-data、import-data、refresh-data、import-group预览、导入与刷新POST /api/connectors/column-values智能筛选列值自动补全POST /api/connectors/connect、disconnect、get-status连接生命周期管理重要约定不要为每个数据源新增专有后端路由或前端面板。除非产品需求已经超出数据源浏览与导入的范围否则优先通过 loader 接口表达能力。数据流从源码看加载链路是一条明确的 Arrow 直通管道见 external_data_loader.py外部数据源 → PyArrow Table → ParquetWorkspace存储格式数据以 parquet 写入 WorkspaceDuckDB 仅作为计算引擎不参与存储。内存格式统一使用 PyArrow 作为中间表示。内置的 PostgreSQL、MySQL、MSSQL 分别使用psycopg2、pymysql、pyodbc读取源数据后再转换为pyarrow.Table。导入入口基类ingest_to_workspace(workspace, ...)自动调用fetch_data_as_arrow()并写入 Workspace 与元数据Loader 无需自己实现 ingest 逻辑。2. 两种接入场景2.1 已有 loader type零代码接入如果系统已经内置对应 loader type例如 PostgreSQL、MySQL、Superset、S3、BigQuery使用者不需要写任何代码直接在界面操作打开 Load Data 页面。点击Add Connection。选择数据源类型。填写 URL、host、database、token、用户名密码等参数。点击Add Connect。连接成功后会出现一个 connector instance 卡片。一个用户可以创建多个同类型连接例如MySQL · prod和MySQL · staging。内置 loader type 的完整清单定义在 data_loader/init.py 的_LOADER_SPECS中每一项为registry_key, module_path, class_name, pip_package_LOADER_SPECS: list[tuple[str, str, str, str]] [ (mysql, data_formulator.data_loader.mysql_data_loader, MySQLDataLoader, pymysql), (mssql, data_formulator.data_loader.mssql_data_loader, MSSQLDataLoader, mssql-python), (postgresql, data_formulator.data_loader.postgresql_data_loader, PostgreSQLDataLoader, psycopg2-binary), (kusto, data_formulator.data_loader.kusto_data_loader, KustoDataLoader, azure-kusto-data), (databricks, data_formulator.data_loader.databricks_data_loader, DatabricksDataLoader, databricks-sql-connector), (s3, data_formulator.data_loader.s3_data_loader, S3DataLoader, boto3), (azure_blob, data_formulator.data_loader.azure_blob_data_loader, AzureBlobDataLoader, azure-storage-blob), (mongodb, data_formulator.data_loader.mongodb_data_loader, MongoDBDataLoader, pymongo), (cosmosdb, data_formulator.data_loader.cosmosdb_data_loader, CosmosDBDataLoader, azure-cosmos), (bigquery, data_formulator.data_loader.bigquery_data_loader, BigQueryDataLoader, google-cloud-bigquery), (athena, data_formulator.data_loader.athena_data_loader, AthenaDataLoader, boto3), (superset, data_formulator.data_loader.superset_data_loader, SupersetLoader, requests), (local_folder, data_formulator.data_loader.local_folder_data_loader, LocalFolderDataLoader, pyarrow), (sample_datasets, data_formulator.data_loader.sample_datasets_loader, SampleDatasetsLoader, requests), ]每个内置 loader 都通过独立的 try/except 导入见 data_loader/init.py缺少可选依赖时只禁用那一个 loader不会让整个应用启动失败并会在DISABLED_LOADERS中给出pip install package提示。2.2 全新的数据源类型外部插件如果要接入一个全新的报表系统或内部数据服务管理员可以提供一个外部 loader 文件~/.data_formulator/plugins/my_report_data_loader.py也可以通过环境变量显式指定插件目录DF_PLUGIN_DIR/opt/data-formulator-loaders服务启动时会扫描插件目录中所有*_data_loader.py文件registry key 由文件名推导my_report_data_loader.py - my_report注册成功后my_report会出现在 Add Connection 的可选数据源类型中。插件目录解析顺序从 data_loader/init.py 的_resolve_plugin_dir()可以看出目录解析优先级DF_PLUGIN_DIR环境变量显式覆盖适合团队共享目录、只读挂载与开发迭代DATA_FORMULATOR_HOME/plugins默认位置与其他 DF 制品保持一致~/.data_formulator/plugins/当DATA_FORMULATOR_HOME未设置时的最终回退。插件扫描的安全门控插件文件会在服务器进程内执行任意 Python 代码因此扫描并非无条件开启。_plugin_scanning_enabled()data_loader/init.py的规则是当WORKSPACE_BACKENDlocal单用户本地模式默认时自动启用多用户 / 托管部署WORKSPACE_BACKEND ! local时必须显式设置DF_ALLOW_PLUGINS1才启用插件目录必须是受信任目录——加载插件等于在服务器进程中执行任意 Python 代码。同时出于防凭据窃取的安全考量插件不允许覆盖内置 loader否则一个恶意的mysql_data_loader.py可能替换内置 MySQL 连接器并窃取所有既有连接的密码重复 key 也会被拒绝这些被拒绝的尝试会以结构化记录进入PLUGIN_ERRORS供 UI 展示见 data_loader/init.py。相关门控与失败路径均由测试覆盖见 test_plugin_scanner.py。说明插件机制是管理员 / 部署方在不改源码的情况下接入新数据源的方式不是普通浏览器终端用户上传任意代码的自助插件机制。3. 最小 Loader 示例以下是一个完整的、可运行的 loader 骨架以接入一个报表系统 API为例即官方文档中的最小示例from typing import Any import pyarrow as pa from data_formulator.data_loader.external_data_loader import ExternalDataLoader class MyReportDataLoader(ExternalDataLoader): def __init__(self, params: dict[str, Any]): self.params params self.url params.get(url, ).rstrip(/) self.token params.get(token, ) if not self.url: raise ValueError(url is required) staticmethod def list_params() - list[dict[str, Any]]: return [ { name: url, type: string, required: True, tier: connection, description: Report system base URL, }, { name: token, type: password, required: True, sensitive: True, tier: auth, description: API token with dataset read access, }, ] staticmethod def auth_instructions() - str: return Provide the report system URL and an API token with dataset read access. def list_tables(self, table_filter: str | None None) - list[dict[str, Any]]: return [ { name: sales_report, metadata: { row_count: None, columns: [{name: region, type: STRING}], sample_rows: [], }, } ] def fetch_data_as_arrow( self, source_table: str, import_options: dict[str, Any] | None None, ) - pa.Table: # Call the source system API and convert the result to Arrow. return pa.table({region: [US, EU]})官方完整示例SQLite 插件仓库 examples/plugins/sqlite_data_loader.py 提供了一个零第三方依赖的完整实战模板覆盖了上述示例之外的更多细节非常值得参考通过DISPLAY_NAME SQLite覆盖 registry key 的默认 Title-case 显示名避免显示成Sqlite在__init__中用modero的 SQLite URI 以只读方式打开数据库并校验文件存在在list_params()中声明default默认值字段在auth_instructions()中返回带 Markdown 格式与示例命令的说明文本在list_tables()中通过sqlite_master列举表与视图并用PRAGMA table_info补充列类型在fetch_data_as_arrow()中实现size受MAX_IMPORT_ROWS上限约束与sort_columns/sort_order的排序支持。快速验证该示例# 1. 单用户模式运行 Data FormulatorWORKSPACE_BACKEND 未设置或为 local # 2. 把文件复制到 ~/.data_formulator/plugins/或 DF_PLUGIN_DIR 指向的目录 # 3. 重启服务Add Connection 中应出现 sqlite sqlite3 /tmp/demo.db SQL CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT, age INTEGER); INSERT INTO users VALUES (1, Alice, 30), (2, Bob, 25); SQL然后在 Data Formulator 中添加一个指向/tmp/demo.db的 SQLite 连接即可。4. 必须实现的方法方法说明__init__(self, params)保存并验证连接参数可创建客户端或延迟到首次调用时连接list_params()返回连接参数定义用于 Add Connection 表单auth_instructions()返回给用户看的连接或认证说明list_tables(table_filterNone)返回可导入表、文件或数据集列表fetch_data_as_arrow(source_table, import_optionsNone)读取数据并返回pyarrow.Table这五个方法是抽象基类abstractmethod中必须实现的契约见 external_data_loader.py 与 L649-L724。其中两个方法的关键语义如下。4.1fetch_data_as_arrow()只接受对象标识符不接受原始 SQL该方法只接受表、文件、集合或数据集标识符source_table不应该接受用户传入的原始 SQL。这是为了同时保证安全性与方言一致性——让前端直接下发 SQL 会引入注入风险也会让 loader 无法屏蔽各数据源间的方言差异。如果支持筛选、排序、列选择或行数限制应从import_options中读取{ size: 100000, columns: [id, amount], sort_columns: [created_at], sort_order: desc, filters: [], source_filters: [], }字段语义结合源码确认size最大行数上限。必须通过min(opts.get(size, MAX_IMPORT_ROWS), MAX_IMPORT_ROWS)截断到硬上限MAX_IMPORT_ROWS 2_000_000200 万行防止意外 OOM见 external_data_loader.py。SQLite 示例正是这样处理的。columns列投影控制返回哪些列。sort_columns/sort_order在行数截断前的排序依据与方向asc/desc。filters标准 SPJ 筛选条件projection、selection、join 意义上的过滤条件集合。source_filters源系统定义如 BI 工具的筛选条件使用EQ/NEQ/GTE/ILIKE等源无关运算符词表合法运算符见下文 §7.2。基类还提供了fetch_data_as_dataframe()便捷方法内部调用fetch_data_as_arrow()后转 pandas性能优先时建议直接实现并返回 Arrow。4.2list_tables()扁平列举与轻量元数据list_tables(table_filterNone)是扁平 / 急切的列举接口返回当前 pinned scope 下所有可导入对象。它与懒加载的ls()长期共存并非遗留接口见 external_data_loader.py。每个返回项包含name用于导入的表标识符如public.usersmetadata含row_count、columns、sample_rows的元数据字典path可选显式层级路径列表如[public, users]存在时用于构建目录树否则按点号拆分name。轻量元数据契约list_tables()用于全量轻量 catalog 与 Agent 搜索缓存必须避免昂贵的逐表数据查询返回表名、稳定源标识符、列名、列类型以及源端低成本可得的表/列描述不要为了 catalog 列表执行 per-tableSELECT * LIMIT ...、COUNT(*)或读取大文件样本row_count和sample_rows仅在源系统已低成本提供时才返回否则省略SQL 类 loader 应优先批量读取information_schema/ 系统 catalog / SDK metadata。推荐的轻量返回形状{ name: public.orders, metadata: { _source_name: public.orders, # 稳定源标识符preview/import/refresh 依赖 description: 订单事实表, columns: [ {name: order_id, type: INTEGER, description: 订单唯一标识}, {name: created_at, type: TIMESTAMP}, ], }, }5. 推荐实现的方法方法用途catalog_hierarchy()声明目录层级例如 database - schema - tablels(pathNone, filterNone, limitNone, offset0)懒加载当前层级节点大目录应支持分页search_catalog(query, limit100)跨层级搜索表或数据集get_metadata(path)返回表/列元数据、描述、样例等get_column_types(source_table)返回源系统列类型帮助前端选择筛选控件get_column_values(source_table, column_name, keyword, limit, offset)返回列值候选项用于智能筛选test_connection()用轻量请求验证连接是否可用5.1 能力对应关系目标能力Loader 接口列举所有表/数据集list_tables()按层级浏览catalog_hierarchy()ls()当前层级过滤ls(..., filter...)跨层级搜索search_catalog(query)表/列元数据get_metadata(path)、get_column_types(source_table)列值枚举get_column_values(...)预览、导入、刷新fetch_data_as_arrow()5.2catalog_hierarchy()与层级浏览catalog_hierarchy()声明从根到叶的完整层级见 external_data_loader.py。每个条目含key内部标识当该层级可通过连接参数固定pinnable时与list_params()中的参数名对应如databaselabel面向用户的显示名如Database。内置数据源的层级声明示例# MySQL [{key: database, label: Database}, {key: table, label: Table}] # PostgreSQL [{key: database, label: Database}, {key: schema, label: Schema}, {key: table, label: Table}] # BigQuery [{key: project, label: Project}, {key: dataset, label: Dataset}, {key: table, label: Table}] # S3 [{key: bucket, label: Bucket}, {key: object, label: File}]作用域固定scope pinning当连接参数匹配某个层级 key例如用户提供了databaseanalytics该层级被固定并从浏览树中隐藏effective_hierarchy()会计算实际可浏览层级。例如 PostgreSQL 提供了databaseprod时完整层级database → schema → table变为可浏览的schema → table。懒加载接口ls()path[]返回第一层展开节点时返回下一层filter只作用于当前层级。对可能返回大量节点的层级例如 database/schema 下有几千张表ls()应支持limit/offset并由/api/connectors/get-catalog透传。get_metadata(path)的默认实现正是通过ls(path[:-1], filterpath[-1])查找节点元数据。search_catalog(query, limit100)大目录应覆盖此方法只返回轻量搜索结果默认实现回退到list_tables(table_filterquery)可能逐表拉取列、样例与计数不适合大目录。前端只在 Enter/搜索按钮时调用/api/connectors/search-catalog。5.3 智能筛选相关方法get_column_types(source_table)返回源系统列类型与描述供 preview/import/refresh 阶段使用返回格式{ description: 订单事实表, columns: [ {name: id, type: NUMERIC, description: 订单 ID, is_dttm: False}, {name: created_at, type: TEMPORAL, description: 创建时间, is_dttm: True}, {name: active, type: BOOLEAN, is_dttm: False}, {name: name, type: STRING, is_dttm: False}, ], }type是源系统原始类型如TIMESTAMP、VARCHAR、BOOLEAN不是 pandas dtype。标准化的source_type分类为TEMPORAL、NUMERIC、BOOLEAN、STRING前端据此选择筛选控件日期选择器、数值范围滑块、布尔开关、文本搜索框。基类默认通过get_metadata派生列信息SQL 类 loader 用information_schema实现get_metadata后即可免费获得此能力。get_column_values(source_table, column_name, keyword, limit, offset)返回列值候选智能筛选自动补全格式{ options: [ {label: Alice, value: Alice}, {label: Bob, value: Bob}, ], has_more: False }参数keyword为前端输入过滤关键词limit为返回数量上限1–200offset为分页偏移。基类默认返回{options: [], has_more: False}表示不支持列值查询前端会回退为自由文本输入。5.4test_connection()用轻量请求验证连接是否可用。基类默认实现调用list_tables(table_filter__ping__)external_data_loader.py子类应覆盖为更廉价的调用如 SQL 的SELECT 1。该验证用于连接创建和 Credential Vault 自动重连校验。5.5 其他可选增强DISPLAY_NAME人类可读的 UI 标签为None时/api/data-loaders会对 registry key 做 Title-case 化覆盖它可避免Sqlite、Bigquery、Mysql之类的尴尬显示。DESCRIPTION一句话连接器用途说明供>connectors: - id: my_report_prod type: my_report name: My Report · prod icon: report params: url: ${MY_REPORT_URL}6.2 环境变量方式也可以用环境变量创建扁平 key 对应嵌套结构DF_SOURCES__my_report_prod__typemy_report DF_SOURCES__my_report_prod__nameMy Report · prod DF_SOURCES__my_report_prod__params__urlhttps://report.example.com6.3 凭据边界重要不要把密码、token、API key、connection string 写入connectors.yaml。敏感凭证只在连接时使用并由 Credential Vault 按用户identity和连接source_id隔离保存。从源码看连接参数分为两类Connector metadatahost、port、database、bucket、root_dir 等非敏感、非 auth-tier 参数可持久化到users/identity/connectors.yaml与Credentials用户名、密码、token、access key、connection string、auth-tier 参数只用于本次连接按identity source_id存入 vault。基类get_safe_params()external_data_loader.py会以list_params()的sensitive标志与type password为准回退到SENSITIVE_PARAMS名称集合password、api_key、secret、token、access_token、refresh_token、access_key、secret_key剔除敏感值只把安全参数写入持久化元数据。7. 认证与凭证7.1 基本凭证认证简单账号密码或 token 认证通常只需要在list_params()中声明敏感字段{name: password, type: password, sensitive: True, tier: auth}tier取值约定connection连接元数据、auth认证专用字段。validate_params()会基于声明校验必填参数并抛出ConnectorParamError当skip_auth_tierTrue时可跳过 auth 层校验适用于 SSO/token 流程。auth_paths()会把tierauth字段自动包装成互斥的认证路径默认credentials。7.2 高级认证模式如果数据源支持 SSO、token exchange 或 delegated login可以按需实现方法说明auth_config()声明 credentials、sso_exchange、delegated、oauth2 等认证模式delegated_login_config()声明弹窗登录 URL 与按钮文案auth_mode()旧兼容接口新 loader 优先使用auth_config()auth_config()支持的认证模式见 external_data_loader.pymodecredentials默认通过 Vault 的静态用户名/密码modesso_exchangeSSO token → 目标系统 token 的后端到后端交换必需exchange_url可选token_url、login_url弹窗回退、timeoutmodedelegated弹窗打开目标系统登录页通过postMessage回传必需login_url可选token_urlmodeoauth2独立 OAuth2 流程不同 IdP必需authorize_url、token_url可选scopes、client_id_env、client_secret_env。以内置 Superset loader 为参考见 superset_data_loader.py它声明了sso_exchange模式并通过PLG_SUPERSET_URL、PLG_SUPERSET_SSO_LOGIN_URL环境变量注入 exchange/login 地址delegated_login_config()返回弹窗登录 URL弹窗预期通过postMessage回传df-sso-auth消息含access_token、refresh_token、user。认证细节见 docs/dev-guides/4-authentication-oidc-tokenstore.md。8. 安全要求不要执行用户提供的原始 SQL。fetch_data_as_arrow()只接受表/文件/数据集标识符。表名、列名、文件路径等外部输入必须校验或转义。例如 SQLite 示例中的_quote_ident()对标识符做双引号转义基类提供build_where_clause()/build_where_clause_inline()/build_source_filter_where_clause_inline()external_data_loader.py等受控 WHERE 构建工具运算符白名单为 ! LIKE NOT LIKE ILIKE IN NOT IN BETWEEN IS NULL IS NOT NULL标识符中若含分号、空字节或--、/*注释序列会被直接拒绝_DANGEROUS_IDENT_RE。涉及本机文件系统读取的 loader 必须限制在安全目录并在多人服务器模式下默认禁用。实现上通过data_loader/__init__.py的_enforce_deployment_restrictions()data_loader/init.py注册禁用规则当WORKSPACE_BACKEND ! local时local_folder这类接受用户可控本机路径的 loader 会自动移入DISABLED_LOADERScreate_connector()会拒绝已禁用类型。list_params()中的密码、token、secret、access key 等字段必须标记为sensitive或typepassword确保get_safe_params()与 vault 正确隔离。返回给前端的元数据不能包含密钥、连接串或内部路径日志与错误信息同样不得写入凭据、token、连接串或敏感路径。auth_config()不包含明文密钥client id、client secret、token URL 等可部署差异通过环境变量或非敏感 params 注入。9. 验证清单开发完成后按以下清单自查loader 文件名是*_data_loader.py。文件中定义了公开的ExternalDataLoader子类。list_params()、auth_instructions()、list_tables()、fetch_data_as_arrow()已实现。大目录数据源实现了ls(..., limit, offset)或search_catalog()。表节点需要稳定源标识符时metadata[_source_name]已设置。每条 table 记录包含table_key字段数据源内稳定唯一标识基类ensure_table_keys()提供 fallback但 loader 应主动设置。list_tables()返回轻量 metadata不做逐表采样、计数或大数据读取。fetch_data_as_arrow()尊重import_options中的size、columns、sort_columns、sort_order、filters、source_filters且size受MAX_IMPORT_ROWS200 万硬上限约束。敏感字段不会写入connectors.yaml也不会返回给前端。外部 API 调用设置了合理 timeout缺少可选依赖时只禁用该 loader不让整个应用启动失败。服务重启后GET /api/data-loaders能看到新的 loader type。Add Connection 能创建连接preview/import/refresh 能正常工作。若 loader 访问用户可控的宿主文件路径已在_enforce_deployment_restrictions()中注册多用户部署禁用规则。10. 深入阅读完整 loader 开发规范含数据流、catalog 浏览 API 契约、部署安全限制、测试要求docs/dev-guides/3-data-loader-development.mdDataConnector API 与连接生命周期docs/dev-guides/5-data-connector-api.md认证与 TokenStoredocs/dev-guides/4-authentication-oidc-tokenstore.md统一行数上限MAX_IMPORT_ROWSdocs/dev-guides/13-unified-row-limits.mdDataFrame 序列化安全df_to_safe_recordsdocs/dev-guides/15-dataframe-serialization.md完整可运行的 SQLite 插件模板examples/plugins/sqlite_data_loader.py插件扫描器测试门控、失败路径、覆盖语义tests/backend/data/test_plugin_scanner.py【免费下载链接】data-formulator Data Formulator is an interactive AI-powered data analysis system makes it easy to connect, explore and visualize data.项目地址: https://gitcode.com/GitHub_Trending/da/data-formulator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考