在 Hadoop 生态使用 JuiceFS
JuiceFS 提供与 HDFS 接口高度兼容的 Java 客户端,Hadoop 生态中的各种应用都可以在不改变代码的情况下,平滑地使用 JuiceFS 存储数据。
环境要求
1. Hadoop 及相关组件
JuiceFS Hadoop Java SDK 同时兼容 Hadoop 2.x、Hadoop 3.x,以及 Hadoop 生态中的各种主流组件。
2. 用户权限
JuiceFS 默认使用本地的「用户/UID」及「用户组/GID」映射,在分布式环境下使用时,为了避免权限问题,请参考文档将需要使用的「用户/UID」及「用户组/GID」同步到所有 Hadoop 节点。也可以通过定义一个全局的用户和用户组文件使得集群中的所有节点共享权限配置,相关配置请查看这里。
3. 文件系统
通过 JuiceFS Java 客户端为 Hadoop 生态提供存储,需要提前创建 JuiceFS 文件系统。部署 Java 客户端时,在配置文件中指定已创建文件系统的元数据引擎地址。
创建文件系统可以参考 JuiceFS 快速上手指南。
如果要在分布式环境中使用 JuiceFS,创建文件系统时,请合理规划要使用的对象存储和数据库,确保它们可以被每个集群节点正常访问。
4. 内存资源
根据计算任务(如 Spark executor)的读写负载,JuiceFS Hadoop Java SDK 可能需要额外使用 4 * juicefs.memory-size 的堆外内存用来加速读写性能。默认情况下,建议为计算任务至少配置 1.2GB 的堆外内存。
5. Java 运行时版本
JuiceFS Hadoop Java SDK 默认使用 JDK 8 编译,如果需要在高版本的 Java 运行时中使用(如 Java 17),需在 JVM 参数中增加以下选项以允许使用反射 API:
--add-exports=java.base/sun.nio.ch=ALL-UNNAMED
更多关于以上选项的说明请参考官方文档。
安装与编译客户端
安装预编译客户端
请参考「安装」文档了解如何下载预编译的 JuiceFS Hadoop Java SDK。
手动编译客户端
不论为哪个系统环境编译客户端,编译后的 JAR 文件都为相同的名称,且只能部署在匹配的系统环境中,例如在 Linux 中编译则只能用于 Linux 环境。另外,由于编译的包依赖 glibc,建议尽量使用低版本的系统进行编译,这样可以获得更好的兼容性。
编译依赖以下工具:
- Go 1.15+(中国用户建议使用 Goproxy China 镜像加速)
- JDK 8+
- Maven 3.3+(中国用户建议使用阿里云镜像加速)
- Git
- make
- GCC 5.4+
Linux 和 macOS
克隆仓库:
git clone https://github.com/juicedata/juicefs.git
进入目录,执行编译:
cd juicefs/sdk/java
make
如果使用 Ceph 的 RADOS 作为 JuiceFS 的存储引擎,需要先安装 librados-dev 包。
cd juicefs/sdk/java
make ceph
编译完成后,可以在 sdk/java/target 目录中找到编译好的 JAR 文件,包括两个版本:
- 包含第三方依赖的包:
juicefs-hadoop-X.Y.Z.jar - 不包含第三方依赖的包:
original-juicefs-hadoop-X.Y.Z.jar
建议使用包含第三方依赖的版本。
Windows
用于 Windows 环境的客户端需要在 Linux 或 macOS 系统上通过交叉编译的方式获得,编译依赖 mingw-w64,需要提前安装。
与编译面向 Linux 和 macOS 客户端的步骤相同,比如在 Ubuntu 系统上,先安装 mingw-w64 包,解决依赖问题:
sudo apt install mingw-w64
克隆并进入 JuiceFS 源代码目录,执行以下代码进行编译:
cd juicefs/sdk/java
make win
部署客户端
让 Hadoop 生态各组件能够正确识别 JuiceFS,需要进行以下配置:
- 将编译好的 JAR 文件和
$JAVA_HOME/lib/tools.jar放置到组件的classpath内,常见大数据平台和组件的安装路径见下表。 - 将 JuiceFS 相关配置写入配置文件(通常是
core-site.xml),详见客户端配置参数。
建议将 JAR 文件放置在一个统一的位置,其他位置通过符号链接进行调用。
大数据平台
| 名称 | 安装路径 |
|---|---|
| CDH | /opt/cloudera/parcels/CDH/lib/hadoop/lib/opt/cloudera/parcels/CDH/spark/jars/var/lib/impala |
| HDP | /usr/hdp/current/hadoop-client/lib/usr/hdp/current/hive-client/auxlib/usr/hdp/current/spark2-client/jars |
| Amazon EMR | /usr/lib/hadoop/lib/usr/lib/spark/jars/usr/lib/hive/auxlib |
| 阿里云 EMR | /opt/apps/ecm/service/hadoop/*/package/hadoop*/share/hadoop/common/lib/opt/apps/ecm/service/spark/*/package/spark*/jars/opt/apps/ecm/service/presto/*/package/presto*/plugin/hive-hadoop2/opt/apps/ecm/service/hive/*/package/apache-hive*/lib/opt/apps/ecm/service/impala/*/package/impala*/lib |
| 腾讯云 EMR | /usr/local/service/hadoop/share/hadoop/common/lib/usr/local/service/presto/plugin/hive-hadoop2/usr/local/service/spark/jars/usr/local/service/hive/auxlib |
| UCloud UHadoop | /home/hadoop/share/hadoop/common/lib/home/hadoop/hive/auxlib/home/hadoop/spark/jars/home/hadoop/presto/plugin/hive-hadoop2 |
| 百度云 EMR | /opt/bmr/hadoop/share/hadoop/common/lib/opt/bmr/hive/auxlib/opt/bmr/spark2/jars |
社区开源组件
| 名称 | 安装路径 |
|---|---|
| Hadoop | ${HADOOP_HOME}/share/hadoop/common/lib/, ${HADOOP_HOME}/share/hadoop/mapreduce/lib/ |
| Spark | ${SPARK_HOME}/jars |
| Presto | ${PRESTO_HOME}/plugin/hive-hadoop2 |
| Trino | ${TRINO_HOME}/plugin/hive |
| Flink | ${FLINK_HOME}/lib |
| StarRocks | ${StarRocks_HOME}/fe/lib/, ${StarRocks_HOME}/be/lib/hadoop/common/lib |
客户端配置参数
请参考以下表格设置 JuiceFS 文件系统相关参数,并写入配置文件,一般是 core-site.xml。
核心配置
| 配置项 | 默认值 | 描述 |
|---|---|---|
fs.jfs.impl | io.juicefs.JuiceFileSystem | 指定要使用的存储实现,默认使用 jfs:// 作为 scheme。如想要使用其它 scheme(例如 cfs://),则修改为 fs.cfs.impl 即可。无论使用的 scheme 是什么,访问的都是 JuiceFS 中的数据。 |
fs.AbstractFileSystem.jfs.impl | io.juicefs.JuiceFS | 指定要使用的存储实现,默认使用 jfs:// 作为 scheme。如想要使用其它 scheme(例如 cfs://),则修改为 fs.AbstractFileSystem.cfs.impl 即可。无论使用的 scheme 是什么,访问的都是 JuiceFS 中的数据。 |
juicefs.meta | 指定预先创建好的 JuiceFS 文件系统的元数据引擎地址。可以通过 juicefs.{vol_name}.meta 格式为客户端同时配置多个文件系统。具体请参考「多文件系统配置」。 |
缓存配 置
| 配置项 | 默认值 | 描述 |
|---|---|---|
juicefs.cache-dir | 设置本地缓存目录,可以指定多个文件夹,用冒号 : 分隔,也可以使用通配符(比如 * )。请预先创建好这些目录,并给予 0777 权限,便于多个应用共享缓存数据。 | |
juicefs.cache-size | 0 | 设置本地缓存目录的容量,单位 MiB,默认为 0,即不开启缓存。如果配置了多个缓存目录,该值代表所有缓存目录容量的总和。 |
juicefs.cache-full-block | true | 是否缓存所有读取的数据块,false 表示只缓存随机读的数据块。 |
juicefs.free-space | 0.1 | 本地缓存目录的最小可用空间比例,默认保留 10% 剩余空间。 |
juicefs.open-cache | 0 | 缓存打开的文件元数据(单位:秒),0 表示关闭 |
juicefs.attr-cache | 0 | 目录和文件属性缓存的过期时间(单位:秒) |
juicefs.entry-cache | 0 | 文件项缓存的过期时间(单位:秒) |
juicefs.dir-entry-cache | 0 | 目录项缓存的过期时间(单位:秒) |
juicefs.discover-nodes-url | 指定发现集群缓存节点列表的方式,每 10 分钟刷新一次。
|
I/O 配置
| 配置项 | 默认值 | 描述 |
|---|---|---|
juicefs.max-uploads | 20 | 上传数据的最大连接数 |
juicefs.max-downloads | 200 | 下载连接的最大数量 |
juicefs.max-deletes | 10 | 删除数据的最大连接数 |
juicefs.get-timeout | 5 | 下载一个对象的超时时间,单位为秒。 |
juicefs.put-timeout | 60 | 上传一个对象的超时时间,单位为秒。 |
juicefs.memory-size | 300 | 读写数据的缓冲区最大空间,单位为 MiB。 |
juicefs.prefetch | 1 | 预读数据块的线程数 |
juicefs.upload-limit | 0 | 上传带宽限制,单位为 Mbps,默认不限制。 |
juicefs.download-limit | 0 | 下载带宽限制,单位为 Mbps,默认不限制。 |
juicefs.io-retries | 10 | IO 失败重试次数 |
juicefs.writeback | false | 是否后台异步上传数据 |
其它配置
| 配置项 | 默认值 | 描述 |
|---|---|---|
juicefs.bucket | 为对象存储指定跟格式化时不同的访问地址 | |
juicefs.debug | false | 是否开启 debug 日志 |
juicefs.access-log | 访问日志的路径。需要所有应用都有写权限,可以配置为 /tmp/juicefs.access.log。该文件会自动轮转,保留最近 7 个文件。 | |
juicefs.superuser | hdfs | 超级用户 |
juicefs.supergroup | supergroup | 超级用户组 |
juicefs.users | null | 用户名以及 UID 列表文件的地址,比如 jfs://name/etc/users。文件格式为 <username>:<UID>,一行一个用户。 |
juicefs.groups | null | 用户组、GID 以及组成员列表文件的地址,比如 jfs://name/etc/groups。文件格式为 <group-name>:<GID>:<username1>,<username2>,一行一个用户组。 |
juicefs.guid-mask | 应用于自动生成 UID/GID 的位掩码,支持非零 uint32 整数。留空时保持现有的按元数据引擎选择掩码的行为。 | |
juicefs.umask | null | 创建文件和目录的 umask 值(如 0022),如果没有此配置,默认值是 fs.permissions.umask-mode。 |
juicefs.push-gateway | Prometheus Pushgateway 地址,格式为 <host>:<port>。 | |
juicefs.push-auth | Prometheus 基本认证信息,格式为 <username>:<password>。 | |
juicefs.push-graphite | Graphite 地址,格式为 <host>:<port>。 | |
juicefs.push-interval | 10 | 指标推送的时间间隔,单位为秒。 |
juicefs.push-labels | 指标额外标签,格式为 key1:value1;key2:value2。 | |
juicefs.fast-resolve | true | 是否开启快速元数据查找(通过 Redis Lua 脚本实现) |
juicefs.no-usage-report | false | 是否上报数据。仅上版本号等使用量数据,不包含任何用户信息。 |
juicefs.block.size | 134217728 | 单位为字节,同 HDFS 的 dfs.blocksize,默认 128 MB |
juicefs.file.checksum | false | DistCp 使用 -update 参数时,是否计算文件 Checksum |
juicefs.no-bgjob | false | 是否关闭后台任务(清理、备份等) |
juicefs.backup-meta | 3600 | 自动将 JuiceFS 元数据备份到对象存储间隔(单位:秒),设置为 0 关闭自动备份 |
juicefs.backup-skip-trash | false | 备份元数据时忽略回收站中的文件和目录。 |
juicefs.heartbeat | 12 | 客户端和元数据引擎之间的心跳间隔(单位:秒),建议所有客户端都设置一样 |
juicefs.skip-dir-mtime | 100ms | 修改父目录 mtime 间隔。 |
juicefs.subdir | 仅允许访问此目录的子路径。可 以指定多个路径,使用逗号分隔。所有其他路径,包括根目录或同级目录,都将被拒绝访问。 |
juicefs.guid-mask 只影响没有显式映射的用户和用户组所生成的 UID/GID,支持任意非零 uint32 整数并自动识别进制:默认为十进制,0b 表示二进制,0 或 0o 表示八进制,0x 表示十六进制。前缀和字母不区分大小写,数字之间可以使用下划线分隔。元数据迁移场景建议使用 0x7fffffff 或 0xffffffff:前者将生成的 ID 限制在有符号 32 位整数的非负范围内,后者保留完整的无符号 32 位值。未配置该参数时,Hadoop SDK 保持原有行为:SQL 元数据引擎使用 0x7fffffff,Redis 和 KV 元数据引擎保留完整的无符号 32 位值。在 MySQL 或 PostgreSQL 元数据引擎上使用大于 0x7fffffff 的值前,需要先将 node 表的 uid 和 gid 字段修改为 BIGINT。
迁移时应选择与源卷匹配的掩码。访问同一个卷的所有 Hadoop SDK 客户端必须使用相同的掩码。启用该参数前,需要将所有 Hadoop SDK 客户端升级到支持 juicefs.guid-mask 的版本;旧版本不会识别该参数,仍会根据当前元数据引擎选择掩码。不要让新旧版本客户端混合访问迁移后的卷。
更换元数据引擎时配置 guid-mask
掩码必须跟随源卷已经使用的 UID/GID 格式,而不是根据目标元数据引擎的类型选择。如果源卷已经显式配置了 guid-mask,目标端应继续使用相同的值。对于此前未配置该参数的卷,按下表配置迁移:
| 迁移方向 | 源卷使用的掩码 | 目标卷配置 | 前置条件 |
|---|---|---|---|
| Redis/KV 转 MySQL | 0xffffffff | juicefs.guid-mask=0xffffffff | 写入迁移的 inode 数据前,MySQL node.uid 和 node.gid 字段必须已经是 BIGINT。 |
| MySQL 转 Redis/KV | 0x7fffffff | juicefs.guid-mask=0x7fffffff | 无。显式配置该掩码可以防止 Redis/KV 客户端改用完整的无符号 32 位 ID。 |
Redis/KV 转 MySQL
将名为 myjfs 的卷从 Redis/KV 迁移到 MySQL 时,所有 Hadoop SDK 客户端都需要配置:
<property>
<name>juicefs.myjfs.guid-mask</name>
<value>0xffffffff</value>
</property>
这个参数只控制 Hadoop SDK 生成 UID/GID 的方式,不会修改 MySQL 表结构。在导入任何 inode 数据前,目标 MySQL 的 node.uid 和 node.gid 字段必须能够保存完整的无符号 32 位值,因此需要使用 BIGINT。
目前标准的 juicefs load 要求目标数据库为空,并会自行创建表后立即导入数据,因此 不能通过“预先创建 node 表并执行 ALTER TABLE”来满足上述条件。如果备份中可能包含大于 2147483647 的 UID/GID,当前版本不支持使用默认表结构执行 Redis/KV 到 MySQL 的标准 juicefs load;必须先提供一种能在建表阶段将这两个字段创建为 BIGINT 的迁移方式。等 inode 数据开始导入后再修改字段,已经无法避免导入过程中的数值越界错误。
MySQL 转 Redis/KV
将 myjfs 从 MySQL 迁移到 Redis/KV 时,配置相同的参数名,并将值设置为:
<property>
<name>juicefs.myjfs.guid-mask</name>
<value>0x7fffffff</value>
</property>
这个方向不需要修改数据库字段。无论迁移方向,迁移前都需要停止所有客户端,迁移完成后再使用对应的 guid-mask 配置重启客户端。
多文件系统配置
当需要同时使用多个 JuiceFS 文件系统时,上述所有配置项均可对特定文件系统进行指定,只需要将文件系统名字放在配置项的中间 ,比如下面示例中的 jfs1 和 jfs2:
<property>
<name>juicefs.jfs1.meta</name>
<value>redis://jfs1.host:port/1</value>
</property>
<property>
<name>juicefs.jfs2.meta</name>
<value>redis://jfs2.host:port/1</value>
</property>
配置示例
以下是一个常用的配置示例,请替换 juicefs.meta 配置中的 {HOST}、{PORT} 和 {DB} 变量为实际的值。
<property>
<name>fs.jfs.impl</name>
<value>io.juicefs.JuiceFileSystem</value>
</property>
<property>
<name>fs.AbstractFileSystem.jfs.impl</name>
<value>io.juicefs.JuiceFS</value>
</property>
<property>
<name>juicefs.meta</name>
<value>redis://{HOST}:{PORT}/{DB}</value>
</property>
<property>
<name>juicefs.cache-dir</name>
<value>/data*/jfs</value>
</property>
<property>
<name>juicefs.cache-size</name>
<value>1024</value>
</property>
<property>
<name>juicefs.access-log</name>
<value>/tmp/juicefs.access.log</value>
</property>
Hadoop 环境配置
请参照前述各项配置表,将配置参数加入到 Hadoop 配置文件 core-site.xml 中。
CDH6
如果使用的是 CDH 6 版本,除了修改 core-site 外,还需要通过 YARN 服务界面修改 mapreduce.application.classpath,增加:
$HADOOP_COMMON_HOME/lib/juicefs-hadoop.jar
HDP
除了修改 core-site 外,还需要通过 MapReduce2 服务界面修改配置 mapreduce.application.classpath,在末尾增加(变量无需替换):
/usr/hdp/${hdp.version}/hadoop/lib/juicefs-hadoop.jar
Flink
将配置参数加入 conf/flink-conf.yaml。如果只是在 Flink 中使用 JuiceFS, 可以不在 Hadoop 环境配置 JuiceFS,只需要配置 Flink 客户端即可。
在阿里云实时平台 Flink SQL 使用 JuiceFS
-
创建 Maven 项目,根据 Flink 不同版本引入如下依赖
<dependencies>
<dependency>
<groupId>io.juicefs</groupId>
<artifactId>juicefs-hadoop</artifactId>
<version>{JUICEFS_HADOOP_VERSION}</version>
</dependency>
<!-- for flink-1.13 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-runtime-blink_2.12</artifactId>
<version>1.13.5</version>
<scope>provided</scope>
</dependency>
<!-- for flink-1.15 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-common</artifactId>
<version>1.15.2</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-files</artifactId>
<version>1.15.2</version>
<scope>provided</scope>
</dependency>
</dependencies> -
创建一个 Java class
public class JuiceFileSystemTableFactory extends FileSystemTableFactory {
@Override
public String factoryIdentifier() {
return "juicefs";
}
} -
Flink table connector 是使用 Java’s Service Provider Interfaces (SPI) 加载自定义实现。 在 resources 按照如下结构创建文件
## for flink-1.13
src/main/resources
├── META-INF
│ └── services
│ └── org.apache.flink.table.factories.Factoryorg.apache.flink.table.factories.Factory文件内容:{YOUR_PACKAGE}.JuiceFileSystemTableFactory -
将填写有 JuiceFS 配置的 core-site.xml 放到 src/main/resources 内:
<configuration>
<property>
<name>fs.juicefs.impl</name>
<value>io.juicefs.JuiceFileSystem</value>
</property>
<property>
<name>juicefs.meta</name>
<value>redis://xxx.redis.rds.aliyuncs.com:6379/0</value>
</property>
...
</configuration>注意由于
jfs://scheme 被阿里其他文件系统占用,所以需要配置fs.juicefs.impl类为 JuiceFS 的实现类,并在后续路径使用juicefs://协议。 -
打包,确保 JAR 内包含 resources 目录下内容
-
通过阿里云实时计算平台控制台->应用->作业开发->connectors 界面上传 JAR 文件
-
测试,将如下 SQL 上线运行,可以在 JuiceFS 的
tmp/tbl目录下发现写入内容CREATE TEMPORARY TABLE datagen_source(
name VARCHAR
) WITH (
'connector' = 'datagen',
'number-of-rows' = '100'
);
CREATE TEMPORARY TABLE jfs_sink (name string)
with (
'connector' = 'juicefs', 'path' = 'juicefs://{VOL_NAME}/tmp/tbl', 'format' = 'csv'
);
INSERT INTO jfs_sink
SELECT
name
from datagen_source;
Hudi
Hudi 自 v0.10.0 版本开始支持 JuiceFS,请确保使用正确的版本。
请参考「Hudi 官方文档」了解如何配置 JuiceFS。
Kafka Connect
可以使用 Kafka Connect 和 HDFS Sink Connector(HDFS 2、HDFS 3)将数据落盘存储到 JuiceFS。
首先需要将 JuiceFS 的 SDK 添加到 Kafka Connect 的 classpath 内,如 /usr/share/java/confluentinc-kafka-connect-hdfs/lib。
在新建 Connect Sink 任务时,做如下配置:
- 指 定
hadoop.conf.dir为包含core-site.xml配置文件的目录,若没有运行在 Hadoop 环境,可创建一个单独目录,如/usr/local/juicefs/hadoop,然后将与 JuiceFS 相关的配置添加到core-site.xml。 - 指定
store.url为以jfs://开头的路径
举例:
# 省略其他配置项...
hadoop.conf.dir=/path/to/hadoop-conf
store.url=jfs://path/to/store
HBase
JuiceFS 适合存储 HBase 的 HFile,但不适合用来保存它的事务日志(WAL),因为将日志持久化到对象存储的时间会远高于持久化到 HDFS 的 DataNode 的内存中。
建议部署一个小的 HDFS 集群来存放 WAL,HFile 文件则存储在 JuiceFS 上。
新建 HBase 集群
修改 hbase-site.xml 配置:
<property>
<name>hbase.rootdir</name>
<value>jfs://{vol_name}/hbase</value>
</property>
<property>
<name>hbase.wal.dir</name>
<value>hdfs://{ns}/hbase-wal</value>
</property>
修改原有 HBase 集群
除了修改上述配置项外,由于 HBase 集群已经在 ZooKeeper 里存储了部分数据,为了避免冲突,有以下两种方式解决:
-
删除原集群
通过 ZooKeeper 客户端删除
zookeeper.znode.parent配置的 znode(默认/hbase)。注意此操作将会删除原有 HBase 上面的所有数据
-
使用新的 znode
保留原 HBase 集群的 znode,以便后续可以恢复。然后为
zookeeper.znode.parent配置一个新的值:hbase-site.xml<property>
<name>zookeeper.znode.parent</name>
<value>/hbase-jfs</value>
</property>
重启服务
当需要使用以下组件访问 JuiceFS 数据时,需要重启相关服务。
在重启之前需要保证 JuiceFS 配置已经写入配置文件,通常可以查看机器上各组件配置的 core-site.xml 里面是否有 JuiceFS 相关配置。
| 组件名 | 服务名 |
|---|---|
| Hive | HiveServer Metastore |
| Spark | ThriftServer |
| Presto | Coordinator Worker |
| Impala | Catalog Server Daemon |
| HBase | Master RegionServer |
HDFS、Hue、ZooKeeper 等服务无需重启。
若访问 JuiceFS 出现 Class io.juicefs.JuiceFileSystem not found 或 No FilesSystem for scheme: jfs 错误,请参考 FAQ。
回收站
JuiceFS Hadoop Java SDK 同样也有和 HDFS 一样的回收站功能,需要通过设置 fs.trash.interval 和 fs.trash.checkpoint.interval 开启,请参考 HDFS 文档了解更多信息。