OpenMetadata SFTP 连接器接入指南:从连接配置到文件元数据采集
【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata
SFTP(SSH File Transfer Protocol)连接器是 OpenMetadata Drive(驱动)类服务家族的一员,用于把 SFTP 服务器上的目录与文件作为可治理的数据资产纳入元数据平台。本文基于仓库中的 SFTP 连接器文档(Sftp.md)与对应源码实现,完整讲解 SFTP 服务的连接配置、认证方式、过滤规则与元数据采集行为,读者读完可以独立完成一个 SFTP 服务的创建、测试连接与元数据管道配置。
一、SFTP 连接器在 OpenMetadata 中的定位
OpenMetadata 通过统一的Drive Service抽象来治理"文件型"数据源(Google Drive、SharePoint、OneDrive、SFTP 等)。从源码的拓扑定义看,drive_service.py 中DriveServiceTopology描述的层级为:
service -> directory -> file / spreadsheet -> worksheetSFTP 是该抽象下的一种具体实现,其服务规范注册于 service_spec.py:
ServiceSpec = BaseSpec(metadata_source_class=SftpSource, connection_class=SftpConnection)即 SFTP 使用SftpSource(采集器)与SftpConnection(连接器)组合。与 Google Drive 不同,SFTP 场景下不涉及 Spreadsheet/Worksheet 实体——在 metadata.py 中,get_spreadsheets_list、yield_spreadsheet、yield_worksheet等方法均返回空,因此 SFTP 采集只产出Directory(目录)与File(文件)两类实体。
二、前置要求
要采集 SFTP 服务器上的元数据,用于登录的账号必须对需要编目的目录和文件具有读取权限。这是文档明确列出的唯一硬性要求——连接器本身只做元数据读取(列出目录、读取文件头、提取 CSV/TSV 结构),不会对服务器上的数据做任何写操作。
从源码角度看,SFTP 底层基于paramiko库实现(见 connection.py 的import paramiko),通过Transport+SFTPClient建立通道;连接的建立、鉴权、资源关闭都在该模块中完成。
三、连接配置详解(Connection Details)
SFTP 服务的连接配置在 UI 中以表单形式呈现,其字段与默认值由 JSON Schema 定义在 sftpConnection.json。该 Schema 要求host与authType为必填项,其余字段均有默认值或可省略。
1. Host(主机)
SFTP 服务器的主机名或 IP 地址,例如sftp.example.com或192.168.1.100。对应 Schema 中host字段,类型为字符串。
2. Port(端口)
SFTP 服务器端口,默认22。源码中同样以connection.port or 22作为兜底逻辑(见 connection.py),即未显式配置时使用 22。如果服务器运行在非标准端口(如企业内部常用的 2222),需在此显式指定。
3. Authentication Type(认证类型)
SFTP 支持两种认证方式,对应 Schema 中authType的oneOf分支:
| 认证方式 | Schema 定义 | 说明 |
|---|---|---|
| Username/Password | basicAuth | 使用用户名 + 密码认证,username与password均为必填 |
| Private Key | keyAuth | 使用 PEM 格式的 SSH 私钥认证,username与privateKey为必填 |
两种方式在源码中的处理路径(connection.py):
- 密码认证:
transport.connect(username=..., password=...) - 私钥认证:先解析私钥为
paramiko.PKey,再transport.connect(username=..., pkey=pkey)
4. Username(用户名)
SFTP 登录用户名,两种认证方式都需要。
5. Password(密码)
密码认证使用的口令。Schema 中该字段标记为format: "password",在 UI 与日志中会作为敏感信息处理。
6. Private Key(私钥)
PEM 格式的 SSH 私钥内容,用于私钥认证。文档与 Schema 描述支持 RSA、Ed25519、ECDSA 与 DSS 四种密钥类型。需要留意的是:从源码实现看,connection.py 的_parse_private_key实际依次尝试paramiko.RSAKey、paramiko.Ed25519Key、paramiko.ECDSAKey三种类型进行解析,且代码注释明确指出paramiko自 4.0.0 起移除了DSSKey(OpenSSH 自 7.0 起也默认拒绝 ssh-dss)。因此在实际使用中,推荐使用 RSA、Ed25519 或 ECDSA 私钥;解析失败时会抛出 "Unable to parse private key. Ensure it is in PEM format" 的错误提示。
7. Private Key Passphrase(私钥口令)
当私钥文件本身被加密(带 passphrase)时,用于解密私钥。若私钥未加密则留空即可。源码中该值作为from_private_key(key_file, password=passphrase)的password参数传入。
8. Root Directories(根目录)
要扫描的文件与子目录根路径列表,默认值为/(即用户家目录)。可以配置多个目录,将采集范围限定在服务器上的特定路径。例如:
/:扫描用户家目录下的全部内容["/data/landing", "/data/export"]:只扫描这两个目录及其子树
在 Schema 中该字段类型为字符串数组,默认["/"]。源码层面,metadata.py 的_fetch_directories会遍历rootDirectories(未配置时取["/"])递归构建目录树,同时记录根目录前缀用于后续 FQN 的路径剥离与重建。
9. Connection Options(连接选项)
额外的连接选项,用于构建发送给服务的连接 URL。对应 Schema 中connectionOptions字段(引用自 connectionBasicType.json 的通用定义)。对 SFTP 而言通常无需配置。
10. Connection Arguments(连接参数)
额外的连接参数,如安全或协议相关的配置,会在连接时发送给服务。同样引用自通用的connectionArguments定义,属于高级可选项。
四、目录与文件过滤
SFTP 连接配置中包含两个过滤模式字段,用于控制采集范围:
- Directory Filter Pattern(目录过滤):正则表达式,只包含/排除匹配的目录。
- File Filter Pattern(文件过滤):正则表达式,只包含/排除匹配的文件。
两者的结构由 filterPattern.json 定义,均为includes与excludes两个正则字符串列表。配置示例:
{ "directoryFilterPattern": { "includes": ["finance.*"], "excludes": ["tmp", "archive.*"] }, "fileFilterPattern": { "includes": [".*\\.csv$", ".*\\.tsv$"], "excludes": [".*\\.log$"] } }过滤逻辑在源码中体现为:
- 目录:
filter_by_directory(self.source_config.directoryFilterPattern, ...)(metadata.pyget_directory_names); - 文件:
filter_by_file(self.source_config.fileFilterPattern, file_info.name)(yield_file方法中,对根目录文件与各目录内文件均执行)。
此外,DriveServiceMetadataPipeline还提供useFqnForFiltering选项(默认false):开启后,正则将作用于完整限定名(如service_name.directory_name.file_name)而非裸名称,适用于需要跨目录精确过滤的场景。
五、结构化数据采集开关
1. Structured Data Files Only(仅结构化数据文件)
默认false。开启后,只编目能够提取出 Schema 的结构化数据文件(CSV、TSV),图片、PDF、视频等非结构化文件会被跳过。
源码中结构化文件的判定集合为(metadata.py):
CSV_EXTENSIONS = {".csv", ".tsv"} STRUCTURED_DATA_EXTENSIONS = {".csv", ".tsv", ".json", ".jsonl", ".parquet", ".avro"}注意区分:_is_csv_file只识别.csv/.tsv(决定是否做 Schema 提取),而_is_structured_data_file使用更宽的扩展名集合(决定是否被采集为 File 实体)。也就是说:
- 关闭该开关时,所有文件都会被编目为 File 实体,但只有 CSV/TSV 会被提取列 Schema;
- 开启该开关时,JSON/Parquet/Avro 等文件也会被编目(它们属于结构化数据),图片/PDF/视频则被跳过。
2. Extract Sample Data(提取样本数据)
默认false,默认关闭以避免性能开销。开启后,会从结构化文件(CSV、TSV)中提取样本数据,写入对应 File 实体,便于在 UI 中预览文件内容。
源码中的样本提取逻辑(_extract_csv_schema与_ingest_sample_data_for_file):
- 远程读取文件内容(
self.client.sftp.open(file_path, "r")),按 UTF-8 解码; - 使用
pandas.read_csv解析,CHUNKSIZE = 200行作为读取上限; - 列类型通过
PANDAS_DTYPE_MAP映射为 OpenMetadata 的DataType:int64/int32 -> INT、float64/float32 -> FLOAT、bool -> BOOLEAN、datetime64[ns] -> DATETIME、object -> STRING,未匹配到的类型统一降级为STRING; - 样本数据取前
MAX_SAMPLE_ROWS = 50行,构建为TableData后通过metadata.ingest_file_sample_data写入。
TSV 文件会自动使用制表符\t作为分隔符(_get_csv_separator按扩展名判断)。
六、采集管道配置(DriveServiceMetadataPipeline)
SFTP 属于 Drive 类服务,其元数据管道类型为DriveMetadata,配置项定义于 driveServiceMetadataPipeline.json。除上述过滤项外,常用配置包括:
| 配置项 | 默认值 | 说明 |
|---|---|---|
includeDirectories | true | 是否采集目录元数据 |
includeFiles | true | 是否采集文件元数据 |
markDeletedDirectories | true | 源端目录被删除后,在 OpenMetadata 中软删除对应目录及其关联实体 |
markDeletedFiles | true | 源端文件被删除后,软删除对应文件 |
includeTags | true | 是否采集标签 |
includeOwners | false | 是否将匹配的负责人关联到实体 |
overrideMetadata | false | 是否用源端元数据覆盖服务器上已有的描述、标签、负责人等 |
threads | 1 | 并行采集线程数 |
markDeletedDirectories/markDeletedFiles的实现位于 drive_service.py 的mark_directories_as_deleted/mark_files_as_deleted,通过delete_entity_from_source对照本次采集的实体状态集做差集删除。SFTP 实现中重写了register_record_directory与register_record_file,确保嵌套目录与根目录文件的 FQN 与入库结果严格一致,避免"刚采集就被当作陈旧实体删除"的问题。
一个完整的 SFTP 元数据管道 YAML 骨架如下(serviceConnection部分即前文配置):
source: type: sftp serviceName: my_sftp_server serviceConnection: config: type: Sftp host: sftp.example.com port: 22 authType: username: data_reader password: ${SFTP_PASSWORD} rootDirectories: - /data/landing structuredDataFilesOnly: false extractSampleData: false directoryFilterPattern: excludes: - tmp fileFilterPattern: includes: - .*\.(csv|tsv|json)$ sourceConfig: config: type: DriveMetadata includeDirectories: true includeFiles: true markDeletedFiles: true threads: 2 sink: type: metadata-rest config: {} workflowConfig: openMetadataServerConfig: hostPort: http://localhost:8585/api authProvider: openmetadata七、连接测试与采集流程
连接测试(Test Connection)
SFTP 连接测试由 connection.py 的test_connection实现,包含两个校验步骤:
- CheckAccess:执行
client.sftp.stat(".")验证能否成功认证并访问服务器; - ListDirectories:遍历每个
rootDirectories,执行client.sftp.listdir(root_dir)验证目录可列出,并统计各目录下的条目数。
任一步骤失败都会抛出SourceConnectionException,UI 上会明确显示是哪一步、哪个目录失败,便于排查权限与路径问题。
元数据采集流程
采集器SftpSource(metadata.py)的完整流程为:
- 建立连接并执行连接测试(
close_on_failure保证失败时释放资源); _fetch_directories递归列出目录树(listdir_attr+stat.S_ISDIR判断),构建SftpDirectoryInfo层级缓存;_fetch_all_files列出各目录(含根目录)下的文件,构建SftpFileInfo(文件名、全路径、大小、修改时间、MIME 类型)缓存;_sort_directories_by_hierarchy按"父目录优先"的 DFS 顺序排序目录;- 对每个目录生成
CreateDirectoryRequest(含parent引用以表达嵌套关系),对每个文件生成CreateFileRequest(关联所属目录、MIME、大小,CSV/TSV 附加列 Schema 与可选样本数据); - 管道结束时依据
markDeleted*配置执行陈旧实体软删除,close()清理缓存并关闭连接。
对应的数据模型(models.py)中,SftpFileInfo与SftpDirectoryInfo均为 pydantic 模型,字段覆盖名称、全路径、大小、修改时间、MIME 类型等,是后续实体构造的基础。
八、补充说明
- 本连接器文档对应的示例数据可见 drives 目录(含服务、目录、文件、电子表格的样例 JSON),可作为理解 Drive 类实体结构的参考;
- 仓库中的 SFTP 连接器遵循 OpenMetadata 通用连接器规范,其服务类型在 driveService.json 中统一注册;
- 配置私钥、口令等敏感信息时,建议通过环境变量或 OpenMetadata 的密钥管理机制注入,避免明文写入管道配置文件。
【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考