
数据工程数据湖大数据对象存储后端【免费下载链接】lakeFSlakeFS - Data version control for your data lake | Git for data项目地址https://gitcode.com/gh_mirrors/la/lakeFS点击查看免费下载lakeFS Hadoop FileSystem下文简称 lakeFSFS是 lakeFS 官方提供的一个org.apache.hadoop.fs.FileSystem实现它让 Spark、Hadoop 等 JVM 大数据组件能够直接以lakefs://协议读写 lakeFS 仓库同时把对象数据data操作直接下沉到底层对象存储仅把元数据metadata操作交给 lakeFS 服务端。读完本文你将掌握lakeFSFS 的架构定位与两种数据访问模式SIMPLE / PRESIGNED的底层原理、全部fs.lakefs.*配置项与默认值、在 Spark 中的完整接入与提交方式以及如何基于仓库自带测试体系进行构建与验证。一、设计定位为什么需要元数据与数据分离的 Hadoop 文件系统官方 README 对该组件给出了精确定位它是org.apache.hadoop.fs.FileSystem的一个实现允许在 lakeFS 上运行 Spark 作业数据操作直接作用于底层存储仅使用 lakeFS 服务端处理元数据操作。这意味着 lakeFSFS 并不会把数据搬进 lakeFS而是数据路径文件的读写直接发生在底层对象存储如 S3由hadoop-aws/ S3A 等客户端完成性能不受 lakeFS 服务端吞吐限制元数据路径文件是否存在、目录结构、对象列表、版本引用分支/提交等语义全部通过 lakeFS OpenAPI/api/v1完成从而获得 lakeFS 的分支、提交、回滚等版本控制能力。从源码看这一分工清晰体现在 LakeFSFileSystem.java 的类结构与 LakeFSClient.java 的 API 客户端封装中LakeFSClient聚合了ObjectsApi、StagingApi、RepositoriesApi、BranchesApi、CommitsApi、ConfigApi、InternalApi七个 lakeFS API 面而文件内容读写则由StorageAccessStrategy的两个实现SimpleStorageAccessStrategy与PresignedStorageAccessStrategy负责。二、路径模型lakefs://repository/ref/pathlakeFSFS 的 URI 遵循lakefs://仓库名/分支或提交引用/路径三段式结构。路径解析逻辑位于 ObjectLocation.javaURI 片段解析字段说明lakefs://scheme协议名默认 scheme 为lakefsexample1repositoryURI 的 host 部分即 lakeFS 仓库名masterrefURI 路径的第一个段可以是分支名、tag 或提交 IDoutput.txtpath分支内的对象路径可为空此时代表分支根目录例如 Spark 示例 spark_with_lakefs.py 中的lakefs://example1/master/output.txt即表示仓库example1的master分支下的output.txt。值得注意的细节FSConfiguration在读取配置时采用scheme 优先、默认 scheme 兜底的策略见 FSConfiguration.java先尝试fs.当前scheme.key找不到再回退到fs.lakefs.keyConstants.DEFAULT_SCHEME即lakefs。因此即便你为集群配置了自定义的 scheme如lakefs2://也依然可以共用一套fs.lakefs.*全局配置。三、构建与发布Maven 工程、Java 8 与两种 JAR3.1 构建前提与产物根据 README 与 pom.xml项目使用Maven构建构建与测试必须安装 Java 8maven-compiler-plugin的source/target均为1.8该模块位于clients/hadoopfs/顶层 Makefile 提供了统一入口make test-hadoopfs等价于cd clients/hadoopfs mvn testhadoop-common与hadoop-aws以provided作用域声明版本由 profile 控制默认2.7.7意味着运行环境中需自行提供 Hadoop 运行时依赖——这与 CHANGELOG 0.1.8 的说明一致该实现可与任意 Hadoop 版本协作只需 classpath 包含hadoop-common与hadoop-aws及其依赖。pom.xml顶部注释明确了两种产物形态产物构建命令用途普通 JARmvn package供更大的程序依赖通常仍需 assembly/shading 后才能直接使用Über JARassemblymvn -Passembly package已做依赖重定位shade可直接放入 Spark 使用3.2 Shading 细节assemblyprofile 使用maven-shade-plugin将所有依赖重定位到io.lakefs.hadoopfs.shade命名空间下见 pom.xml重定位列表包括org.apache.httpcomponents→io.lakefs.hadoopfs.shade.org.apache.httpcomponentsokio/okhttp3→io.lakefs.hadoopfs.shade.okio/io.lakefs.hadoop.shade.okhttp3com.google.gson、io.gsonfire→io.lakefs.hadoop.shade.gson/io.lakefs.hadoop.shade.gsonfireio.lakefs.clients.sdk→io.lakefs.hadoop.shade.sdk这保证了打进 Spark 后不会与 Spark 自身携带的 HTTP 库、JSON 库发生类冲突。3.3 发布途径README 提到发布需遵循 lakeFS 客户端的发布检查清单具体发布操作在pom.xml中也有明确记录发布到 Maven Centralmvn deploy或mvn -Passembly deploy上传到 S3mvn package s3-storage-wagon:s3-uploadupload-jar-Passembly前缀为 assembly 版由s3-storage-wagon插件上传至treeverse-clients-us-east桶的hadoop/目录见 pom.xml。依赖方面需特别留意0.2.0 起使用io.lakefs:sdk而非io.lakefs:api-client当前版本锁定 SDK 1.67.0。如果你不使用 assembly Über JAR 而自行装配依赖必须同步更换坐标否则无法编译运行。四、在 Spark 中启用 lakeFS 文件系统4.1 注册 FileSystem 实现让 Spark 认识lakefs://scheme只需在 Hadoop 配置中注册实现类LakeFSFileSystem 类注释与 run-test.py 均给出同一套配置spark.hadoop.fs.lakefs.implio.lakefs.LakeFSFileSystem spark.hadoop.fs.lakefs.endpointhttp://lakefs:8000/api/v1 spark.hadoop.fs.lakefs.access.keyAKIAIOSFODNN7EXAMPLE spark.hadoop.fs.lakefs.secret.keywJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY其中endpoint默认值为http://localhost:8000/api/v1见 Constants.javaaccess.key/secret.key即 lakeFS 的访问密钥对。以spark.hadoop.前缀包裹后Spark 会把它们透传给 HadoopConfiguration。4.2 PySpark 读写示例仓库在 spark_with_lakefs.py 给出了最小可运行示例from pyspark.sql import SparkSession spark SparkSession.builder.appName(test_app).getOrCreate() spark._jsc.hadoopConfiguration().set(fs.lakefs.access.key, AKIAIOSFODNN7EXAMPLE) spark._jsc.hadoopConfiguration().set(fs.lakefs.secret.key, wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY) # 通过 lakefs 文件系统写入 samples sc.parallelize([ (abonsantofakemail.com, Alberto, Bonsanto), (dbonsantofakemail.com, Dakota, Bonsanto) ]) samples.saveAsTextFile(lakefs://example1/master/output.txt) # 通过 lakefs 文件系统读取 lines sc.textFile(lakefs://example1/master/input.txt)4.3 三种接入模式的提交配置对比仓库的 Spark 集成测试 run-test.py 覆盖了三种模式可作为生产提交模板模式scheme额外关键配置s3_gateway经 S3 网关s3afs.s3a.endpoint指向 lakeFS 网关、关闭 SSLhadoopfsSIMPLE 直连lakefs追加fs.s3a.access.key/fs.s3a.secret.key/fs.s3a.region供底层 S3A 访问物理存储hadoopfs_presigned预签名lakefs追加spark.hadoop.fs.lakefs.access.modepresigned提交方式二选一见 run-test.py# 方式一使用发布到 Maven 的 assembly 包 spark-submit --packages io.lakefs:hadoop-lakefs-assembly:版本 ... # 方式二使用本地构建的 JAR spark-submit --jars /target/client.jar ...五、配置参数全表含默认值以下配置项全部来自 Constants.java 及 LakeFSClient.java 的读取逻辑。除fs.lakefs.impl等 Spark 注册项外配置键均以fs.lakefs.为前缀在 Spark 中写为spark.hadoop.fs.lakefs.*。5.1 连接与认证配置项默认值说明fs.lakefs.endpointhttp://localhost:8000/api/v1lakeFS OpenAPI 地址结尾多余的/会被自动去除fs.lakefs.access.key无lakeFS 访问密钥basic_auth 模式必需fs.lakefs.secret.key无lakeFS 密钥basic_auth 模式必需fs.lakefs.auth.providerbasic_auth认证方式设为 token provider 类全名时走 JWT 换取流程fs.lakefs.session_id无可选会作为sessionIdCookie 附加到 API 请求5.2 访问模式与列举配置项默认值说明fs.lakefs.access.modeSIMPLESIMPLE直连底层存储或PRESIGNED预签名 URLBetafs.lakefs.list.amount1000每次listObjects请求拉取的条目数即分页块大小fs.lakefs.delete.bulk_size1000递归删除时每个deleteObjects批次的对象数BulkDeleter默认值fs.lakefs.delete.object.no.tombstonefalse实验特性允许的话以不产生 tombstone 的方式删除对象5.3 SDK 超时0.17.0 引入配置项默认值说明fs.lakefs.api.connect.timeout.ms10000连接超时10 秒fs.lakefs.api.read.timeout.ms30000读超时30 秒fs.lakefs.api.write.timeout.ms30000写超时30 秒5.4 API 重试0.18.0 引入配置项默认值说明fs.lakefs.api.retry.max-retries5最大重试次数设为0或负数可关闭重试fs.lakefs.api.retry.initial-backoff.ms1000首次退避基数1 秒fs.lakefs.api.retry.max-backoff.ms20000退避上限20 秒fs.lakefs.api.retry.jitter-factor0.25抖动因子实际退避 指数退避 × (1 ± jitter)5.5 Token Provider 相关配置项默认值说明fs.lakefs.token.duration_seconds无服务端决定换取到的 lakeFS token 的过期时长fs.lakefs.token.aws.access.key/secret.key/session.token无临时 AWS 凭证三件套TemporaryAWSCredentialsLakeFSTokenProvider必需fs.lakefs.token.aws.sts.duration_seconds60每次生成的 STS identity token 有效期60 秒fs.lakefs.token.aws.sts.endpoint无STS 端点如sts.amazonaws.com该 Provider 必需fs.lakefs.token.sts.additional_headersX-Lakefs-Server-IDlakeFS服务端host以key1:value1,key2:value2格式传给 STS 签名的额外头5.6 实验特性commit 压缩配置项默认值说明experimental.commit-every.num-deletes0关闭每累计删除多少对象后按概率触发一次内部 commit辅助 tombstone 压缩experimental.commit-every.prob0.1触发概率建议约等于1/执行器数上述两项来自 0.16.0 的 compaction 实验见 Constants.java实现于 LakeFSFileSystem.java 的addDeletedMaybeCompact当删除计数超过阈值时以概率发起一次内部提交并将递归删除按CONTIGUOUS_DELETION_FACTOR 3加权因为连续 tombstone 对性能伤害更大。5.7 故障排查工具配置配置项默认值说明fs.lakefs.tracer.working_dir无FileSystemTracer 对应的 S3 位置桶名或桶内绝对路径fs.lakefs.tracer.use_lakefs_outputfalse是否返回 lakeFS 侧的调用结果默认返回 S3A 侧结果六、两种数据访问模式SIMPLE 与 PRESIGNED访问模式由fs.lakefs.access.mode决定初始化逻辑在 LakeFSFileSystem.java。两种模式分别对应SimpleStorageAccessStrategy与PresignedStorageAccessStrategy共同实现StorageAccessStrategy接口。6.1 SIMPLE 模式直连底层文件系统写路径见 SimpleStorageAccessStrategy.java调用StagingApi.getPhysicalAddress(..., presignfalse)获取物理地址形如s3://bucket/...PhysicalAddressTranslator.translate()将物理地址转换为 HadoopPath——先校验其匹配StorageConfig返回的blockstoreNamespaceValidityRegex再按 blockstore 类型映射为s3a://URI当前仅支持s3azure分支留有 TODO见 PhysicalAddressTranslator.java在底层文件系统如 S3AFileSystem上创建物理文件并用LinkOnCloseOutputStream包裹关闭流时触发链接MetadataClient.getObjectMetadata()先从FileStatus.getETag()反射取 ETag 与长度失败则回退到AmazonS3Client.getObjectMetadata()见 MetadataClient.java随后LakeFSLinker.link()调用linkPhysicalAddress把 ETag、字节数写入 lakeFS 元数据形成未提交对象见 LakeFSLinker.java。非 overwrite 场景会携带If-None-Match: *防止覆盖已存在对象。读路径statObject(presignfalse)取物理地址 → translate → 直接physicalFs.open(physicalPath, bufSize)读取。6.2 PRESIGNED 模式全程走预签名 URLBeta写路径见 PresignedStorageAccessStrategy.javagetPhysicalAddress(..., presigntrue)拿到可 PUT 的预签名 URL通过HttpURLConnection发起PUTContent-Type: application/octet-stream配合LakeFSFileSystemOutputStream在关闭时完成链接。读路径statObject(presigntrue)拿到预签名 GET URL交给 HttpRangeInputStream.java。它先发Range: bytes0-0从Content-Range响应头解析对象总长度随后按 1 MB 缓冲以Range分块拉取并支持seek定位。由于是 HTTP Range 读取不会把整个对象下载到内存。适用场景PRESIGNED 模式无需给集群配置底层存储的凭证如fs.s3a.access.key由 lakeFS 服务端签发带权限的临时 URL是安全、云无关的接入方式代价是读写均经预签名 HTTP 请求完成。6.3 目录标记directory marker语义与 S3A 一致lakeFSFS 用零字节对象 尾随/表示空目录见 LakeFSFileSystem.java 的mkdirs/createDirectoryMarker。写入对象后会通过deleteEmptyDirectoryMarkers向上清理不再需要的标记删除目录后会调用createDirectoryMarkerIfNotExists补回父目录标记。getFileStatus则依次尝试精确 stat →path /标记 stat → 前缀列举三段判定LakeFSFileSystem.java。七、身份认证机制从 Basic Auth 到 AWS STS 换取 lakeFS Token认证入口在 LakeFSClient.javabasic_auth默认直接用fs.lakefs.access.key/fs.lakefs.secret.key配置 HTTP Basic 认证Token Provider当fs.lakefs.auth.provider指向一个类名时LakeFSTokenProviderFactory会通过Class.forName(...)反射加载该实现构造签名为(String scheme, Configuration conf)见 LakeFSTokenProviderFactory.java随后用 JWT Bearer 认证。LakeFSTokenProvider接口只定义两个方法getToken()与refresh()见 LakeFSTokenProvider.java。内置实现有两个TemporaryAWSCredentialsLakeFSTokenProvider.java使用配置中显式给出的临时 AWS 凭证AK/SK/SessionTokenAWSLakeFSTokenProvider.java其换取流程为——先用 AWS 凭证对 STSGetCallerIdentity请求做 SigV4 预签名GetCallerIdentityV4Presigner生成极短生命周期的身份令牌并 Base64 编码再调用 lakeFSAuthApi.externalPrincipalLogin()将其兑换为 lakeFS JWTtoken 缓存在内存中过期前自动刷新AWSLakeFSTokenProvider.java。默认会额外签署X-Lakefs-Server-ID头值为 lakeFS 服务端 host防止请求被重放。这套设计从 CHANGELOG 0.2.4 引入new Token Provider feature with IAM Role Support与 Hadoop-AWS 的自定义凭证提供器思路一致便于对接 EMR、Databricks 等托管环境。八、API 调用的超时与重试保障LakeFSClient.newApiClientNoAuth()LakeFSClient.java在初始化 SDK 客户端时统一应用连接/读/写三档超时5.3 节RetryInterceptorOkHttp 拦截器0.18.0 引入对HTTP 408、429、500~504以及网络IOException进行指数退避重试退避计算为min(initialBackoff × 2^(attempt-1), maxBackoff) × (1 ± jitterFactor)见 RetryInterceptor.java额外附加X-Lakefs-Client: lakefs-hadoopfs/版本请求头方便服务端识别客户端来源。九、核心操作语义与限制9.1 创建createcreate()先检查路径状态目标是目录则抛FileAlreadyExistsException不允许 overwrite 时对已存在文件同样拒绝LakeFSFileSystem.java。文件流关闭前数据只存在于底层存储的 staging 区关闭后才被 lakeFS 看到。9.2 读取open / listFilesopen委托给存储策略创建输入流listFiles返回的LocatedFileStatus通过getFileBlockLocations填充块位置。列目录由内部类ListingIterator驱动以after游标 amount即list.amount分页调用listObjects递归列举时delimiter置空非递归时用/分隔符LakeFSFileSystem.java。9.3 重命名rename——非原子、仅限未提交数据源码注释明确列出了语义LakeFSFileSystem.java仅支持同一分支上的未提交数据跨分支/跨仓库返回falseonSameBranch校验file → existing file覆盖目标file → existing dir移入目录file → 不存在的 dst返回false与 S3A 行为一致directory → existing dir把目录整体移入实现为先 copyObject 后 deleteObject两步非原子中途失败可能留下中间态CHANGELOG 0.1.11 起优先用 CopyObject 替代 StageObject源或目标为根目录时拒绝mtime不保留。9.4 删除delete非空目录必须recursivetrue否则抛IOException递归删除使用BulkDeleter批量调用deleteObjectsAPI单线程执行器、默认批大小 1000见 BulkDeleter.java0.1.7 起借此加速递归删除fs.lakefs.delete.object.no.tombstone开启时删除请求会附带 no-tombstone 标记实验特性。9.5 不支持的操作append直接抛出UnsupportedOperationExceptionAppend is not supported by LakeFSFileSystem——这是对象存储语义与 lakeFS 不可变对象模型共同决定的边界。9.6 块大小与物理地址默认块大小 32 MBDEFAULT_BLOCK_SIZE 32 * ONE_MB见 Constants.java若能从仓库元数据解析出底层存储则委托底层文件系统的默认块大小getFSForConfigLakeFSFileSystem.java。LakeFSFileStatus额外携带checksumETag、physicalAddress、isEmptyDirectory三个扩展字段LakeFSFileStatus.java。十、测试体系单元、契约与集成三层验证10.1 单元测试与 SDK 契约单元测试位于 src/test/java/io/lakefs覆盖FSConfiguration、RetryInterceptor、AWSLakeFSTokenProvider含 SigV4 预签名、LakeFSTokenProviderFactory等并引入了 mockserver 提升速度与精度CHANGELOG 0.2.0。10.2 Hadoop FileSystem 契约测试pom.xml提供四组 profile把 Hadoop 官方契约测试套件跑在 lakeFSFS 之上ProfileHadoop 版本说明contract-tests-hadoop22.7.7默认版本contract-tests-hadoop33.1.4Hadoop 3contract-tests-hadoop3423.4.2Hadoop 3.4.2附加 AWS SDK bundlecontract-tests-hadoop3-presigned3.1.4以lakefs.access_modepresigned运行契约用例定义在 src/test/java/io/lakefs/contract覆盖TestLakeFSFileSystemContract{Create,Delete,Mkdir,Open,Rename,Seek}等标准场景Hadoop 2/3 各有对应子类TestLakeFSFileSystemContractHadoop2/3。运行方式cd clients/hadoopfs mvn -Pcontract-tests-hadoop3 test10.3 本地端到端集成环境test/lakefsfs_contract/docker-compose.yaml 会拉起PostgreSQL元数据库 MinIOS3 兼容对象存储 lakeFS 服务端setup-test.sh 完成初始化创建测试用户、S3 桶并以lakectl repo create创建测试仓库。客户端配置集中在 auth-keys.xmlfs.lakefs.impl、fs.lakefs.endpointhttp://10.5.0.55:8000/api/v1、S3A 凭证等由 core-site.xml 引入。这为复现Spark lakeFSFS S3完整链路提供了开箱即用的基准环境。10.4 Spark 集成测试test/spark/run-test.py 以 Sonnet 词频统计为例通过docker compose拉起 Spark master/worker 与 lakeFS验证s3_gateway、hadoopfs、hadoopfs_presigned三种模式的读写闭环覆盖从上传数据 → spark-submit 提交 → 读取 lakefs 输入 → 写回 lakefs 输出的完整流程。十一、故障排查工具FileSystemTracer当 LakeFSFS 与直接访问底层对象存储的结果出现差异时可启用仓库自带的 FileSystemTracer.java 作为诊断手段把fs.lakefs.impl指向io.lakefs.FileSystemTracer它会在每次文件系统操作时同时调用 LakeFSFileSystem 与 S3AFileSystem对比两者结果并记录[RESULTS_COMPARISON]日志默认返回 S3A 侧输出可用fs.lakefs.tracer.use_lakefs_outputtrue切换。前提是lakefs://repo/branch/与s3a://${fs.lakefs.tracer.working_dir}/下内容一致。这非常适合定位元数据视图与实际对象存储内容不同步类问题。十二、版本演进要点clients/hadoopfs/CHANGELOG.md版本关键变化0.18.0lakeFS API 调用支持可配置的指数退避重试0.17.0支持 SDK 超时配置connect/read/write实验性deleteObjectno-tombstoneSDK 升至 1.67.00.16.0实验周期性 commit 压缩以提升 Spark 性能0.2.5修复根路径通配符listStatus支持 storage ID0.2.4新增 Token ProviderIAM Role 支持0.2.0破坏性变更仅支持 lakeFS 服务端 ≥ 0.108.0依赖从api-client改为sdk0.1.13PRESIGNED 模式Beta0.1.10全部 shade 到io.lakefs.hadoopfs.shadeHadoop 改为 provided 依赖兼容更多 Spark 环境0.1.8可与任意 Hadoop 版本配合classpath 需含 hadoop-common 与 hadoop-aws0.1.7用deleteObjects加速递归删除十三、继续深入源码导航文件系统主体实现LakeFSFileSystem.java配置键与默认值全集Constants.javaAPI 客户端、超时与重试装配LakeFSClient.java两种数据访问策略SimpleStorageAccessStrategy.java、PresignedStorageAccessStrategy.java物理地址转换PhysicalAddressTranslator.javaURI 解析ObjectLocation.java构建与发布pom.xml示例spark_with_lakefs.py使用前提提醒当前实现要求 lakeFS 服务端版本 ≥ 0.108.0对应客户端 0.2.0 的破坏性变更构建与测试需 Java 8SIMPLE 模式的物理地址翻译当前仅支持s3块存储类型append不受支持。这些边界均以本仓库当前代码为准。赞分享数据工程数据湖大数据对象存储后端【免费下载链接】lakeFSlakeFS - Data version control for your data lake | Git for data项目地址https://gitcode.com/gh_mirrors/la/lakeFS点击查看免费下载相关推荐lakeFS 数据版本控制实战以 Git 方式管理数据湖lakeFS 数据版本控制实战以 Git 方式管理数据湖 导读 lakeFS 是一个 source available源代码可用的数据版本控制工具它把数据工程数据湖大数据对象存储后端lakeFS数据湖版本控制的开源利器lakeFS数据湖版本控制的开源利器 lakeFS 是一个开源项目旨在将对象存储转变为类似 Git 的版本控制仓库使得数据湖的管理方式与代码管理方式相类似数据工程数据湖大数据对象存储后端为什么iPad检测不到smartbanner.js平台检测原理揭秘maxTouchPoints与User Agent的较量为什么iPad检测不到smartbanner.js平台检测原理揭秘maxTouchPoints与User Agent的较量 如果你发现 smartbanne上一篇终极指南如何在LMStudio中使用革命性的Qwen3-235B-A22B-Thinking-2507-FP8模型下一篇突破数据实时写入瓶颈StarRocks Stream Load事务接口深度解析与实践指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考