diff --git a/.github/workflows/nightly-ci.yml b/.github/workflows/nightly-ci.yml
index a948dbcc20f..f7f955233a1 100644
--- a/.github/workflows/nightly-ci.yml
+++ b/.github/workflows/nightly-ci.yml
@@ -102,7 +102,7 @@ jobs:
working-directory: tests-integration/fixtures
run: ../../.github/scripts/pull-test-deps-images.sh && docker compose up -d --wait
- name: Run nextest cases
- run: cargo nextest run --workspace -F dashboard -F pg_kvbackend -F mysql_kvbackend -F ai_functions
+ run: cargo nextest run --workspace -F dashboard -F pg_kvbackend -F mysql_kvbackend -F ai_functions -F hdfs-object-store
env:
CARGO_BUILD_RUSTFLAGS: "-C link-arg=-fuse-ld=mold"
RUST_BACKTRACE: 1
diff --git a/.github/workflows/rust.yml b/.github/workflows/rust.yml
index d4ab7e1d7b3..e29b69678c1 100644
--- a/.github/workflows/rust.yml
+++ b/.github/workflows/rust.yml
@@ -233,7 +233,7 @@ jobs:
run: ../../.github/scripts/pull-test-deps-images.sh && docker compose up -d --wait
- name: Run nextest cases
- run: cargo nextest run --workspace -F dashboard -F pg_kvbackend -F mysql_kvbackend -F ai_functions
+ run: cargo nextest run --workspace -F dashboard -F pg_kvbackend -F mysql_kvbackend -F ai_functions -F hdfs-object-store
env:
CARGO_BUILD_RUSTFLAGS: "-C link-arg=-fuse-ld=mold"
RUST_BACKTRACE: 1
@@ -290,7 +290,7 @@ jobs:
run: ../../.github/scripts/pull-test-deps-images.sh && docker compose up -d --wait
- name: Run nextest cases
- run: cargo llvm-cov nextest --workspace --lcov --output-path lcov.info -F dashboard -F pg_kvbackend -F mysql_kvbackend -F ai_functions
+ run: cargo llvm-cov nextest --workspace --lcov --output-path lcov.info -F dashboard -F pg_kvbackend -F mysql_kvbackend -F ai_functions -F hdfs-object-store
env:
CARGO_BUILD_RUSTFLAGS: "-C link-arg=-fuse-ld=mold"
RUST_BACKTRACE: 1
diff --git a/Cargo.lock b/Cargo.lock
index bb7f69d1551..4ced42fc860 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -40,10 +40,21 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b169f7a6d4742236a0a00c541b845991d0ac43e546831af1249753ab4c3aa3a0"
dependencies = [
"cfg-if",
- "cipher",
+ "cipher 0.4.4",
"cpufeatures 0.2.17",
]
+[[package]]
+name = "aes"
+version = "0.9.3"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "35f0f96ce78e38c3dc6d8948aa8163d06385be74000f3c7a95bf1eef35d3ea32"
+dependencies = [
+ "cipher 0.5.2",
+ "cpubits",
+ "cpufeatures 0.3.0",
+]
+
[[package]]
name = "aes-siv"
version = "0.7.0"
@@ -51,10 +62,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7e08d0cdb774acd1e4dac11478b1a0c0d203134b2aab0ba25eb430de9b18f8b9"
dependencies = [
"aead",
- "aes",
- "cipher",
+ "aes 0.8.4",
+ "cipher 0.4.4",
"cmac",
- "ctr",
+ "ctr 0.9.2",
"dbl",
"digest 0.10.7",
"zeroize",
@@ -1400,6 +1411,15 @@ dependencies = [
"generic-array",
]
+[[package]]
+name = "block-padding"
+version = "0.4.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "710f1dd022ef4e93f8a438b4ba958de7f64308434fa6a87104481645cc30068b"
+dependencies = [
+ "hybrid-array",
+]
+
[[package]]
name = "blocking"
version = "1.6.1"
@@ -1534,9 +1554,9 @@ checksum = "40e38929add23cdf8a366df9b0e088953150724bcbe5fc330b0d8eb3b328eec8"
[[package]]
name = "bumpalo"
-version = "3.19.0"
+version = "3.20.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "46c5e41b57b8bba42a04676d81cb89e9ee8e859a1a66f80a5a72e1cb76b34d43"
+checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649"
[[package]]
name = "bytecheck"
@@ -1730,7 +1750,16 @@ version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "26b52a9543ae338f279b96b0b9fed9c8093744685043739079ce85cd58f289a6"
dependencies = [
- "cipher",
+ "cipher 0.4.4",
+]
+
+[[package]]
+name = "cbc"
+version = "0.2.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "ce2dc9ee5f88d11e0beb842c88b33c8a5cf0d1329c4b19494af42b07dbfe8896"
+dependencies = [
+ "cipher 0.5.2",
]
[[package]]
@@ -1775,7 +1804,7 @@ version = "0.8.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "738b8d467867f80a71351933f70461f5b56f24d5c93e0cf216e59229c968d330"
dependencies = [
- "cipher",
+ "cipher 0.4.4",
]
[[package]]
@@ -1827,7 +1856,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c3613f74bd2eac03dad61bd53dbe620703d4371614fe0bc3b9f04dd36fe4e818"
dependencies = [
"cfg-if",
- "cipher",
+ "cipher 0.4.4",
"cpufeatures 0.2.17",
]
@@ -1850,7 +1879,7 @@ checksum = "10cd79432192d1c0f4e1a0fef9527696cc039165d729fb41b3f4f4f354c2dc35"
dependencies = [
"aead",
"chacha20 0.9.1",
- "cipher",
+ "cipher 0.4.4",
"poly1305",
"zeroize",
]
@@ -1943,10 +1972,21 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "773f3b9af64447d2ce9850330c473515014aa235e6a783b02db81ff39e4a3dad"
dependencies = [
"crypto-common 0.1.6",
- "inout",
+ "inout 0.1.4",
"zeroize",
]
+[[package]]
+name = "cipher"
+version = "0.5.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "e8cf2a2c93cd704877c0858356ed03480ff301ee950b43f1cbe4573b088bfa6c"
+dependencies = [
+ "block-buffer 0.12.0",
+ "crypto-common 0.2.2",
+ "inout 0.2.2",
+]
+
[[package]]
name = "clang-sys"
version = "1.8.1"
@@ -1955,7 +1995,7 @@ checksum = "0b023947811758c97c59bf9d1c188fd619ad4718dcaa767947df1cadb14f39f4"
dependencies = [
"glob",
"libc",
- "libloading",
+ "libloading 0.8.8",
]
[[package]]
@@ -2118,7 +2158,7 @@ version = "0.7.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8543454e3c3f5126effff9cd44d562af4e31fb8ce1cc0d3dcd8f084515dbc1aa"
dependencies = [
- "cipher",
+ "cipher 0.4.4",
"dbl",
"digest 0.10.7",
]
@@ -2228,7 +2268,7 @@ checksum = "fe6d2e5af09e8c8ad56c969f2157a3d4238cebc7c55f0a517728c38f7b200f81"
dependencies = [
"serde",
"termcolor",
- "unicode-width 0.1.14",
+ "unicode-width 0.2.1",
]
[[package]]
@@ -3218,6 +3258,12 @@ dependencies = [
"cfg-if",
]
+[[package]]
+name = "cpubits"
+version = "0.1.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "15b85f9c39137c3a891689859392b1bd49812121d0d61c9caf00d46ed5ce06ae"
+
[[package]]
name = "cpufeatures"
version = "0.2.17"
@@ -3238,9 +3284,9 @@ dependencies = [
[[package]]
name = "crc"
-version = "3.3.0"
+version = "3.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "9710d3b3739c2e349eb44fe848ad0b7c8cb1e42bd87ee49371df2f7acaf3e675"
+checksum = "5eb8a2a1cd12ab0d987a5d5e825195d372001a4094a0376319d5a0ad71c1ba0d"
dependencies = [
"crc-catalog",
]
@@ -3438,7 +3484,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b9d6cf87adf719ddf43a805e92c6870a531aedda35ff640442cbaf8674e141e1"
dependencies = [
"aead",
- "cipher",
+ "cipher 0.4.4",
"generic-array",
"poly1305",
"salsa20",
@@ -3483,7 +3529,16 @@ version = "0.9.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0369ee1ad671834580515889b80f2ea915f23b8be8d0daa4bbaf2ac5c7590835"
dependencies = [
- "cipher",
+ "cipher 0.4.4",
+]
+
+[[package]]
+name = "ctr"
+version = "0.10.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "baaca1c4b237092596f64d571e9db6ce4109c4ef9742e27590f1709594461f21"
+dependencies = [
+ "cipher 0.5.2",
]
[[package]]
@@ -4742,6 +4797,15 @@ dependencies = [
"syn 2.0.117",
]
+[[package]]
+name = "des"
+version = "0.9.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "916a94e407b54f9034d71dd748234cd1e516ced6284009906ae246f177eafe5a"
+dependencies = [
+ "cipher 0.5.2",
+]
+
[[package]]
name = "diff"
version = "0.1.13"
@@ -5819,6 +5883,34 @@ dependencies = [
"byteorder",
]
+[[package]]
+name = "g2gen"
+version = "1.2.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "c5a7e0eb46f83a20260b850117d204366674e85d3a908d90865c78df9a6b1dfc"
+dependencies = [
+ "g2poly",
+ "proc-macro2",
+ "quote",
+ "syn 2.0.117",
+]
+
+[[package]]
+name = "g2p"
+version = "1.2.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "539e2644c030d3bf4cd208cb842d2ce2f80e82e6e8472390bcef83ceba0d80ad"
+dependencies = [
+ "g2gen",
+ "g2poly",
+]
+
+[[package]]
+name = "g2poly"
+version = "1.2.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "312d2295c7302019c395cfb90dacd00a82a2eabd700429bba9c7a3f38dbbe11b"
+
[[package]]
name = "generator"
version = "0.8.5"
@@ -6167,6 +6259,47 @@ dependencies = [
"hashbrown 0.15.4",
]
+[[package]]
+name = "hdfs-native"
+version = "0.14.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "8fd181084003308224efddf737417832186839ce2882d2468ddad0114cfb3551"
+dependencies = [
+ "aes 0.9.3",
+ "base64 0.22.1",
+ "bitflags 2.12.1",
+ "bumpalo",
+ "bytes",
+ "cbc 0.2.1",
+ "chrono",
+ "cipher 0.5.2",
+ "crc",
+ "ctr 0.10.1",
+ "des",
+ "dns-lookup",
+ "futures",
+ "g2p",
+ "hex",
+ "hmac 0.13.0",
+ "libc",
+ "libloading 0.9.0",
+ "log",
+ "md-5 0.11.0",
+ "num-traits",
+ "once_cell",
+ "prost 0.14.1",
+ "prost-types 0.14.1",
+ "rand 0.10.1",
+ "regex",
+ "roxmltree",
+ "socket2 0.6.4",
+ "thiserror 2.0.17",
+ "tokio",
+ "url",
+ "uuid",
+ "whoami 2.1.3",
+]
+
[[package]]
name = "hdrhistogram"
version = "7.5.4"
@@ -6549,7 +6682,7 @@ dependencies = [
"libc",
"percent-encoding",
"pin-project-lite",
- "socket2 0.5.10",
+ "socket2 0.6.4",
"tokio",
"tower-service",
"tracing",
@@ -6620,7 +6753,7 @@ dependencies = [
"js-sys",
"log",
"wasm-bindgen",
- "windows-core 0.57.0",
+ "windows-core 0.61.2",
]
[[package]]
@@ -6996,10 +7129,20 @@ version = "0.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "879f10e63c20629ecabbb64a8010319738c66a5cd0c29b02d63d272b03751d01"
dependencies = [
- "block-padding",
+ "block-padding 0.3.3",
"generic-array",
]
+[[package]]
+name = "inout"
+version = "0.2.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "4250ce6452e92010fdf7268ccc5d14faa80bb12fc741938534c58f16804e03c7"
+dependencies = [
+ "block-padding 0.4.2",
+ "hybrid-array",
+]
+
[[package]]
name = "instant"
version = "0.1.13"
@@ -7061,7 +7204,7 @@ version = "0.9.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "96e4f67dbfc0f75d7b65953ecf0be3fd84ee0cb1ae72a00a4aa9a2f5518a2c80"
dependencies = [
- "aes",
+ "aes 0.8.4",
]
[[package]]
@@ -7810,6 +7953,16 @@ dependencies = [
"windows-targets 0.53.5",
]
+[[package]]
+name = "libloading"
+version = "0.9.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "754ca22de805bb5744484a5b151a9e1a8e837d5dc232c2d7d8c2e3492edc8b60"
+dependencies = [
+ "cfg-if",
+ "windows-link 0.2.1",
+]
+
[[package]]
name = "liblzma"
version = "0.4.6"
@@ -9300,7 +9453,7 @@ version = "0.7.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "77e878c846a8abae00dd069496dbe8751b16ac1c3d6bd2a7283a938e8228f90d"
dependencies = [
- "proc-macro-crate 1.3.1",
+ "proc-macro-crate 3.3.0",
"proc-macro2",
"quote",
"syn 2.0.117",
@@ -9407,6 +9560,7 @@ dependencies = [
"common-test-util",
"derive_builder 0.20.2",
"futures",
+ "hdfs-native",
"humantime-serde",
"object-store",
"object_store",
@@ -9482,7 +9636,7 @@ version = "0.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2cc40678e045ff4eb1666ea6c0f994b133c31f673c09aed292261b6d5b6963a0"
dependencies = [
- "cipher",
+ "cipher 0.4.4",
]
[[package]]
@@ -9555,6 +9709,7 @@ dependencies = [
"opendal-service-azblob",
"opendal-service-fs",
"opendal-service-gcs",
+ "opendal-service-hdfs-native",
"opendal-service-http",
"opendal-service-mysql",
"opendal-service-oss",
@@ -9745,6 +9900,20 @@ dependencies = [
"uuid",
]
+[[package]]
+name = "opendal-service-hdfs-native"
+version = "0.59.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "1b84340e11a3bcd7f1f0b55df032d040b2907483c546c3021485aba749c3072c"
+dependencies = [
+ "bytes",
+ "futures",
+ "hdfs-native",
+ "log",
+ "opendal-core",
+ "serde",
+]
+
[[package]]
name = "opendal-service-http"
version = "0.59.2"
@@ -10792,8 +10961,8 @@ version = "0.7.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e847e2c91a18bfa887dd028ec33f2fe6f25db77db3619024764914affe8b69a6"
dependencies = [
- "aes",
- "cbc",
+ "aes 0.8.4",
+ "cbc 0.1.2",
"der",
"pbkdf2",
"scrypt",
@@ -11313,7 +11482,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "22505a5c94da8e3b7c2996394d1c933236c4d743e81a410bcca4e6989fc066a4"
dependencies = [
"bytes",
- "heck 0.4.1",
+ "heck 0.5.0",
"itertools 0.12.1",
"log",
"multimap",
@@ -11333,8 +11502,8 @@ version = "0.14.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ac6c3320f9abac597dcbc668774ef006702672474aad53c6d596b62e487b40b1"
dependencies = [
- "heck 0.4.1",
- "itertools 0.10.5",
+ "heck 0.5.0",
+ "itertools 0.14.0",
"log",
"multimap",
"once_cell",
@@ -11382,7 +11551,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d"
dependencies = [
"anyhow",
- "itertools 0.10.5",
+ "itertools 0.14.0",
"proc-macro2",
"quote",
"syn 2.0.117",
@@ -11395,7 +11564,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9120690fafc389a67ba3803df527d0ec9cbbc9cc45e4cc20b332996dfb672425"
dependencies = [
"anyhow",
- "itertools 0.10.5",
+ "itertools 0.14.0",
"proc-macro2",
"quote",
"syn 2.0.117",
@@ -12988,7 +13157,7 @@ version = "0.10.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "97a22f5af31f73a954c10289c93e8a50cc23d971e80ee446f1f6f7137a088213"
dependencies = [
- "cipher",
+ "cipher 0.4.4",
]
[[package]]
@@ -13778,7 +13947,7 @@ version = "0.8.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1961e2ef424c1424204d3a5d6975f934f56b6d50ff5732382d84ebf460e147f7"
dependencies = [
- "heck 0.4.1",
+ "heck 0.5.0",
"proc-macro2",
"quote",
"syn 2.0.117",
@@ -16119,13 +16288,13 @@ version = "0.33.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "772206835b5bca8719825d2ab9e6fb160cbfce8ac1dc0e5dffe22cbe77254468"
dependencies = [
- "aes",
+ "aes 0.8.4",
"aes-siv",
"base16",
"base62",
"base64-simd",
"bytes",
- "cbc",
+ "cbc 0.1.2",
"cfb-mode",
"cfg-if",
"chacha20poly1305",
@@ -16141,7 +16310,7 @@ dependencies = [
"crc",
"crypto_secretbox",
"csv",
- "ctr",
+ "ctr 0.9.2",
"digest 0.10.7",
"dns-lookup",
"domain",
diff --git a/Cargo.toml b/Cargo.toml
index e057db65826..5f497a3faff 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -161,6 +161,7 @@ fst = "0.4.7"
futures = "0.3"
futures-util = "0.3"
greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "9533078d6fd4bdaddbca0750cb9bf7ad64b2e568" }
+hdfs-native = "0.14.6"
hex = "0.4"
hostname = "0.4.0"
http = "1"
diff --git a/config/config.md b/config/config.md
index e380f075548..cebc5091b4c 100644
--- a/config/config.md
+++ b/config/config.md
@@ -148,9 +148,11 @@
| `storage` | -- | -- | The data storage options. |
| `storage.data_home` | String | `./greptimedb_data` | The working home directory. |
| `storage.copy_root` | String | `./greptimedb_data/copy` | Root directory for standalone SQL access to local files.
Relative SQL paths are resolved below this directory. Absolute paths are accepted only when
they are inside this directory. Defaults to `/copy`.
Distributed deployments always reject SQL access to local files.
Upgrade note: COPY commands and existing external tables that reference paths outside this
directory will fail. Move those files below the copy root, set this option to a dedicated
directory containing them, or migrate the files to object storage before upgrading. |
-| `storage.type` | String | `File` | The storage type used to store the data.
- `File`: the data is stored in the local file system.
- `S3`: the data is stored in the S3 object storage.
- `Gcs`: the data is stored in the Google Cloud Storage.
- `Azblob`: the data is stored in the Azure Blob Storage.
- `Oss`: the data is stored in the Aliyun OSS. |
+| `storage.type` | String | `File` | The storage type used to store the data.
- `File`: the data is stored in the local file system.
- `S3`: the data is stored in the S3 object storage.
- `Gcs`: the data is stored in the Google Cloud Storage.
- `Azblob`: the data is stored in the Azure Blob Storage.
- `Oss`: the data is stored in the Aliyun OSS.
- `Hdfs`: the data is stored in the Hadoop Distributed File System. |
| `storage.bucket` | String | Unset | The S3 bucket name.
**It's only used when the storage type is `S3`, `Oss` and `Gcs`**. |
-| `storage.root` | String | Unset | The S3 data will be stored in the specified prefix, for example, `s3://${bucket}/${root}`.
**It's only used when the storage type is `S3`, `Oss` and `Azblob`**. |
+| `storage.root` | String | Unset | The directory or object prefix under which data is stored.
**It's only used when the storage type is `S3`, `Oss`, `Gcs`, `Azblob` and `Hdfs`**. |
+| `storage.name_node` | String | Unset | The HDFS NameNode URI, for example, `hdfs://127.0.0.1:9000`.
**It's only used when the storage type is `Hdfs`**. |
+| `storage.options` | InlineTable | Unset | Additional options passed to the native HDFS client.
**It's only used when the storage type is `Hdfs`**. |
| `storage.access_key_id` | String | Unset | The access key id of the aws account.
It's **highly recommended** to use AWS IAM roles instead of hardcoding the access key id and secret key.
**It's only used when the storage type is `S3` and `Oss`**. |
| `storage.secret_access_key` | String | Unset | The secret access key of the aws account.
It's **highly recommended** to use AWS IAM roles instead of hardcoding the access key id and secret key.
**It's only used when the storage type is `S3`**. |
| `storage.access_key_secret` | String | Unset | The secret access key of the aliyun account.
**It's only used when the storage type is `Oss`**. |
@@ -609,9 +611,11 @@
| `query.experimental_spill_compression` | String | `uncompressed` | Compression for spilled data files: "uncompressed" (default), "lz4_frame", "zstd".
Ignored unless mode is "custom". |
| `storage` | -- | -- | The data storage options. |
| `storage.data_home` | String | `./greptimedb_data` | The working home directory. |
-| `storage.type` | String | `File` | The storage type used to store the data.
- `File`: the data is stored in the local file system.
- `S3`: the data is stored in the S3 object storage.
- `Gcs`: the data is stored in the Google Cloud Storage.
- `Azblob`: the data is stored in the Azure Blob Storage.
- `Oss`: the data is stored in the Aliyun OSS. |
+| `storage.type` | String | `File` | The storage type used to store the data.
- `File`: the data is stored in the local file system.
- `S3`: the data is stored in the S3 object storage.
- `Gcs`: the data is stored in the Google Cloud Storage.
- `Azblob`: the data is stored in the Azure Blob Storage.
- `Oss`: the data is stored in the Aliyun OSS.
- `Hdfs`: the data is stored in the Hadoop Distributed File System. |
| `storage.bucket` | String | Unset | The S3 bucket name.
**It's only used when the storage type is `S3`, `Oss` and `Gcs`**. |
-| `storage.root` | String | Unset | The S3 data will be stored in the specified prefix, for example, `s3://${bucket}/${root}`.
**It's only used when the storage type is `S3`, `Oss` and `Azblob`**. |
+| `storage.root` | String | Unset | The directory or object prefix under which data is stored.
**It's only used when the storage type is `S3`, `Oss`, `Gcs`, `Azblob` and `Hdfs`**. |
+| `storage.name_node` | String | Unset | The HDFS NameNode URI, for example, `hdfs://127.0.0.1:9000`.
**It's only used when the storage type is `Hdfs`**. |
+| `storage.options` | InlineTable | Unset | Additional options passed to the native HDFS client.
**It's only used when the storage type is `Hdfs`**. |
| `storage.access_key_id` | String | Unset | The access key id of the aws account.
It's **highly recommended** to use AWS IAM roles instead of hardcoding the access key id and secret key.
**It's only used when the storage type is `S3` and `Oss`**. |
| `storage.secret_access_key` | String | Unset | The secret access key of the aws account.
It's **highly recommended** to use AWS IAM roles instead of hardcoding the access key id and secret key.
**It's only used when the storage type is `S3`**. |
| `storage.access_key_secret` | String | Unset | The secret access key of the aliyun account.
**It's only used when the storage type is `Oss`**. |
diff --git a/config/datanode.example.toml b/config/datanode.example.toml
index da2d353fef8..253cc3f2a57 100644
--- a/config/datanode.example.toml
+++ b/config/datanode.example.toml
@@ -300,6 +300,13 @@ overwrite_entry_start_id = false
# credential = "base64-credential"
# endpoint = "https://storage.googleapis.com"
+# Example of using HDFS as the storage with the native Rust client.
+# [storage]
+# type = "Hdfs"
+# root = "/greptimedb"
+# name_node = "hdfs://127.0.0.1:9000"
+# options = { "dfs.client.block.write.replace-datanode-on-failure.enable" = "true" }
+
## The query engine options.
[query]
## Parallelism of the query engine.
@@ -348,6 +355,7 @@ data_home = "./greptimedb_data"
## - `Gcs`: the data is stored in the Google Cloud Storage.
## - `Azblob`: the data is stored in the Azure Blob Storage.
## - `Oss`: the data is stored in the Aliyun OSS.
+## - `Hdfs`: the data is stored in the Hadoop Distributed File System.
type = "File"
## The S3 bucket name.
@@ -355,11 +363,21 @@ type = "File"
## @toml2docs:none-default
bucket = "greptimedb"
-## The S3 data will be stored in the specified prefix, for example, `s3://${bucket}/${root}`.
-## **It's only used when the storage type is `S3`, `Oss` and `Azblob`**.
+## The directory or object prefix under which data is stored.
+## **It's only used when the storage type is `S3`, `Oss`, `Gcs`, `Azblob` and `Hdfs`**.
## @toml2docs:none-default
root = "greptimedb"
+## The HDFS NameNode URI, for example, `hdfs://127.0.0.1:9000`.
+## **It's only used when the storage type is `Hdfs`**.
+## @toml2docs:none-default
+name_node = "hdfs://127.0.0.1:9000"
+
+## Additional options passed to the native HDFS client.
+## **It's only used when the storage type is `Hdfs`**.
+## @toml2docs:none-default
+options = {}
+
## The access key id of the aws account.
## It's **highly recommended** to use AWS IAM roles instead of hardcoding the access key id and secret key.
## **It's only used when the storage type is `S3` and `Oss`**.
diff --git a/config/standalone.example.toml b/config/standalone.example.toml
index e0a8deb9c22..fad4f64f08a 100644
--- a/config/standalone.example.toml
+++ b/config/standalone.example.toml
@@ -509,6 +509,13 @@ max_running_procedures = 128
# credential = "base64-credential"
# endpoint = "https://storage.googleapis.com"
+# Example of using HDFS as the storage with the native Rust client.
+# [storage]
+# type = "Hdfs"
+# root = "/greptimedb"
+# name_node = "hdfs://127.0.0.1:9000"
+# options = { "dfs.client.block.write.replace-datanode-on-failure.enable" = "true" }
+
## The query engine options.
[query]
## Parallelism of the query engine.
@@ -566,6 +573,7 @@ data_home = "./greptimedb_data"
## - `Gcs`: the data is stored in the Google Cloud Storage.
## - `Azblob`: the data is stored in the Azure Blob Storage.
## - `Oss`: the data is stored in the Aliyun OSS.
+## - `Hdfs`: the data is stored in the Hadoop Distributed File System.
type = "File"
## The S3 bucket name.
@@ -573,11 +581,21 @@ type = "File"
## @toml2docs:none-default
bucket = "greptimedb"
-## The S3 data will be stored in the specified prefix, for example, `s3://${bucket}/${root}`.
-## **It's only used when the storage type is `S3`, `Oss` and `Azblob`**.
+## The directory or object prefix under which data is stored.
+## **It's only used when the storage type is `S3`, `Oss`, `Gcs`, `Azblob` and `Hdfs`**.
## @toml2docs:none-default
root = "greptimedb"
+## The HDFS NameNode URI, for example, `hdfs://127.0.0.1:9000`.
+## **It's only used when the storage type is `Hdfs`**.
+## @toml2docs:none-default
+name_node = "hdfs://127.0.0.1:9000"
+
+## Additional options passed to the native HDFS client.
+## **It's only used when the storage type is `Hdfs`**.
+## @toml2docs:none-default
+options = {}
+
## The access key id of the aws account.
## It's **highly recommended** to use AWS IAM roles instead of hardcoding the access key id and secret key.
## **It's only used when the storage type is `S3` and `Oss`**.
diff --git a/src/cmd/Cargo.toml b/src/cmd/Cargo.toml
index a67dd895999..2c925de57c8 100644
--- a/src/cmd/Cargo.toml
+++ b/src/cmd/Cargo.toml
@@ -22,6 +22,7 @@ required-features = ["dev-tools"]
[features]
default = [
"ai_functions",
+ "hdfs-object-store",
"servers/pprof",
"servers/mem-prof",
"meta-srv/pg_kvbackend",
@@ -32,6 +33,7 @@ enterprise = ["common-meta/enterprise", "frontend/enterprise", "meta-srv/enterpr
# Developer-only helper binaries and diagnostic datanode commands.
# Kept out of `default` so normal/release builds don't compile them.
dev-tools = []
+hdfs-object-store = ["mito2/hdfs-object-store", "object-store/hdfs-object-store"]
mysql-object-store = ["object-store/mysql-object-store"]
tokio-console = ["common-telemetry/tokio-console"]
diff --git a/src/cmd/tests/load_config_test.rs b/src/cmd/tests/load_config_test.rs
index d4b19812df2..cc095c5c92e 100644
--- a/src/cmd/tests/load_config_test.rs
+++ b/src/cmd/tests/load_config_test.rs
@@ -125,6 +125,44 @@ fn test_load_runtime_options_without_max_blocking_threads() {
);
}
+#[test]
+#[cfg(feature = "hdfs-object-store")]
+fn test_load_datanode_hdfs_config() {
+ let config = tempfile::NamedTempFile::new().unwrap();
+ std::fs::write(
+ config.path(),
+ r#"
+ [storage]
+ type = "Hdfs"
+ name = "local-hdfs"
+ root = "/greptimedb"
+ name_node = "hdfs://127.0.0.1:9000"
+ enable_read_cache = false
+ options = { "dfs.client.block.write.replace-datanode-on-failure.enable" = "true" }
+ "#,
+ )
+ .unwrap();
+
+ let options =
+ GreptimeOptions::::load_layered_options(config.path().to_str(), "")
+ .unwrap();
+ let object_store::config::ObjectStoreConfig::Hdfs(hdfs) = options.component.storage.store
+ else {
+ unreachable!()
+ };
+
+ assert_eq!("local-hdfs", hdfs.name);
+ assert_eq!("/greptimedb", hdfs.connection.root);
+ assert_eq!("hdfs://127.0.0.1:9000", hdfs.connection.name_node);
+ assert_eq!(
+ Some(&"true".to_string()),
+ hdfs.connection
+ .options
+ .get("dfs.client.block.write.replace-datanode-on-failure.enable")
+ );
+ assert!(!hdfs.cache.enable_read_cache);
+}
+
#[allow(deprecated)]
#[test]
fn test_load_datanode_example_config() {
diff --git a/src/mito2/Cargo.toml b/src/mito2/Cargo.toml
index 9ab2dee10df..46032deef38 100644
--- a/src/mito2/Cargo.toml
+++ b/src/mito2/Cargo.toml
@@ -10,6 +10,7 @@ test = ["common-test-util", "rstest", "rstest_reuse", "rskafka"]
testing = ["test"]
test-shared-fs-region-migration = []
enterprise = []
+hdfs-object-store = ["object-store/hdfs-object-store"]
[lints]
workspace = true
diff --git a/src/mito2/src/region/utils.rs b/src/mito2/src/region/utils.rs
index 6917315757a..7dfaad72aa7 100644
--- a/src/mito2/src/region/utils.rs
+++ b/src/mito2/src/region/utils.rs
@@ -265,7 +265,20 @@ impl RegionFileCopier {
#[cfg(test)]
mod tests {
+ #[cfg(feature = "hdfs-object-store")]
+ use object_store::ObjectStore;
+ #[cfg(feature = "hdfs-object-store")]
+ use object_store::layers::HdfsCompatibilityLayer;
+ #[cfg(feature = "hdfs-object-store")]
+ use object_store::services::Fs;
+
use super::*;
+ #[cfg(feature = "hdfs-object-store")]
+ use crate::access_layer::AccessLayer;
+ #[cfg(feature = "hdfs-object-store")]
+ use crate::sst::index::intermediate::IntermediateManager;
+ #[cfg(feature = "hdfs-object-store")]
+ use crate::sst::index::puffin_manager::PuffinManagerFactory;
#[test]
fn test_build_copy_file_paths() {
@@ -344,4 +357,54 @@ mod tests {
format!("/table_dir/1_0000000002/index/{}.1.puffin", file_id)
);
}
+
+ #[cfg(feature = "hdfs-object-store")]
+ #[tokio::test]
+ async fn test_copy_region_files_with_hdfs_fallback() {
+ let (temp_dir, puffin_manager) =
+ PuffinManagerFactory::new_for_test_async("hdfs-copy-region").await;
+ let intermediate_manager = IntermediateManager::init_fs(temp_dir.path().to_string_lossy())
+ .await
+ .unwrap();
+ let storage_dir = temp_dir.path().join("storage");
+ std::fs::create_dir(&storage_dir).unwrap();
+ let object_store = ObjectStore::new(Fs::default().root(storage_dir.to_str().unwrap()))
+ .unwrap()
+ .layer(HdfsCompatibilityLayer::new_for_test());
+ let access_layer = Arc::new(AccessLayer::new(
+ "table_dir",
+ PathType::Bare,
+ object_store.clone(),
+ puffin_manager,
+ intermediate_manager,
+ ));
+ let copier = RegionFileCopier::new(access_layer);
+ let source_region_id = RegionId::new(1, 1);
+ let target_region_id = RegionId::new(1, 2);
+ let file_id = FileId::random();
+ let descriptor = FileDescriptor::Data { file_id, size: 8 };
+ let (source_path, target_path) = build_copy_file_paths(
+ source_region_id,
+ target_region_id,
+ descriptor,
+ "table_dir",
+ PathType::Bare,
+ );
+ object_store.write(&source_path, "contents").await.unwrap();
+
+ copier
+ .copy_files(source_region_id, target_region_id, vec![descriptor], 1)
+ .await
+ .unwrap();
+
+ assert_eq!(
+ b"contents",
+ object_store
+ .read(&target_path)
+ .await
+ .unwrap()
+ .to_bytes()
+ .as_ref()
+ );
+ }
}
diff --git a/src/object-store/Cargo.toml b/src/object-store/Cargo.toml
index 4083271a967..15e2069e3bd 100644
--- a/src/object-store/Cargo.toml
+++ b/src/object-store/Cargo.toml
@@ -8,6 +8,7 @@ license.workspace = true
workspace = true
[features]
+hdfs-object-store = ["dep:hdfs-native", "opendal/services-hdfs-native", "uuid"]
mysql-object-store = ["opendal/services-mysql"]
services-memory = ["opendal/services-memory"]
testing = ["derive_builder", "uuid"]
@@ -21,6 +22,7 @@ common-macro.workspace = true
common-runtime.workspace = true
common-telemetry.workspace = true
derive_builder = { workspace = true, optional = true }
+hdfs-native = { workspace = true, optional = true }
humantime-serde.workspace = true
opendal = { version = "0.59.2", features = [
"layers-tracing",
diff --git a/src/object-store/src/config.rs b/src/object-store/src/config.rs
index 1308db92fab..2917bdbe817 100644
--- a/src/object-store/src/config.rs
+++ b/src/object-store/src/config.rs
@@ -12,10 +12,14 @@
// See the License for the specific language governing permissions and
// limitations under the License.
+#[cfg(feature = "hdfs-object-store")]
+use std::collections::HashMap;
use std::time::Duration;
use common_base::readable_size::ReadableSize;
use common_base::secrets::{ExposeSecret, SecretString};
+#[cfg(feature = "hdfs-object-store")]
+use opendal::services::HdfsNative;
#[cfg(feature = "mysql-object-store")]
use opendal::services::Mysql;
use opendal::services::{Azblob, Gcs, Oss, S3};
@@ -34,6 +38,8 @@ pub enum ObjectStoreConfig {
Oss(OssConfig),
Azblob(AzblobConfig),
Gcs(GcsConfig),
+ #[cfg(feature = "hdfs-object-store")]
+ Hdfs(HdfsConfig),
#[cfg(feature = "mysql-object-store")]
Mysql(MysqlConfig),
}
@@ -53,6 +59,8 @@ impl ObjectStoreConfig {
Self::Oss(_) => "Oss",
Self::Azblob(_) => "Azblob",
Self::Gcs(_) => "Gcs",
+ #[cfg(feature = "hdfs-object-store")]
+ Self::Hdfs(_) => "Hdfs",
#[cfg(feature = "mysql-object-store")]
Self::Mysql(_) => "Mysql",
}
@@ -72,6 +80,8 @@ impl ObjectStoreConfig {
Self::Oss(oss) => &oss.name,
Self::Azblob(az) => &az.name,
Self::Gcs(gcs) => &gcs.name,
+ #[cfg(feature = "hdfs-object-store")]
+ Self::Hdfs(hdfs) => &hdfs.name,
#[cfg(feature = "mysql-object-store")]
Self::Mysql(mysql) => &mysql.name,
};
@@ -91,6 +101,8 @@ impl ObjectStoreConfig {
Self::Oss(oss) => Some(&oss.cache),
Self::Azblob(az) => Some(&az.cache),
Self::Gcs(gcs) => Some(&gcs.cache),
+ #[cfg(feature = "hdfs-object-store")]
+ Self::Hdfs(hdfs) => Some(&hdfs.cache),
#[cfg(feature = "mysql-object-store")]
Self::Mysql(mysql) => Some(&mysql.cache),
}
@@ -104,6 +116,8 @@ impl ObjectStoreConfig {
Self::Oss(oss) => Some(&mut oss.cache),
Self::Azblob(az) => Some(&mut az.cache),
Self::Gcs(gcs) => Some(&mut gcs.cache),
+ #[cfg(feature = "hdfs-object-store")]
+ Self::Hdfs(hdfs) => Some(&mut hdfs.cache),
#[cfg(feature = "mysql-object-store")]
Self::Mysql(mysql) => Some(&mut mysql.cache),
}
@@ -289,6 +303,42 @@ impl From<&GcsConnection> for Gcs {
}
}
+/// Connection options for a Hadoop Distributed File System backend.
+#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
+#[serde(default)]
+#[cfg(feature = "hdfs-object-store")]
+pub struct HdfsConnection {
+ /// Working directory for all object-store operations.
+ pub root: String,
+ /// HDFS NameNode URI, for example `hdfs://127.0.0.1:9000`.
+ pub name_node: String,
+ /// Additional options passed to the native HDFS client.
+ pub options: HashMap,
+}
+
+#[cfg(feature = "hdfs-object-store")]
+impl From<&HdfsConnection> for HdfsNative {
+ fn from(connection: &HdfsConnection) -> Self {
+ let root = util::normalize_dir(&connection.root);
+ HdfsNative::default()
+ .root(&root)
+ .name_node(&connection.name_node)
+ .options(connection.options.clone())
+ }
+}
+
+/// Hadoop Distributed File System object storage configuration.
+#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
+#[serde(default)]
+#[cfg(feature = "hdfs-object-store")]
+pub struct HdfsConfig {
+ pub name: String,
+ #[serde(flatten)]
+ pub connection: HdfsConnection,
+ #[serde(flatten)]
+ pub cache: ObjectStorageCacheConfig,
+}
+
#[cfg(feature = "mysql-object-store")]
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
#[serde(default)]
@@ -411,6 +461,20 @@ mod tests {
assert_eq!("test", s3_config.config_name());
assert_eq!("S3", s3_config.provider_name());
+ #[cfg(feature = "hdfs-object-store")]
+ {
+ let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig::default());
+ assert_eq!("Hdfs", hdfs_config.config_name());
+ assert_eq!("Hdfs", hdfs_config.provider_name());
+
+ let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig {
+ name: "test".to_string(),
+ ..Default::default()
+ });
+ assert_eq!("test", hdfs_config.config_name());
+ assert_eq!("Hdfs", hdfs_config.provider_name());
+ }
+
#[cfg(feature = "mysql-object-store")]
{
let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig::default());
@@ -438,6 +502,11 @@ mod tests {
assert!(gcs_config.is_object_storage());
let azblob_config = ObjectStoreConfig::Azblob(AzblobConfig::default());
assert!(azblob_config.is_object_storage());
+ #[cfg(feature = "hdfs-object-store")]
+ {
+ let hdfs_config = ObjectStoreConfig::Hdfs(HdfsConfig::default());
+ assert!(hdfs_config.is_object_storage());
+ }
#[cfg(feature = "mysql-object-store")]
{
let mysql_config = ObjectStoreConfig::Mysql(MysqlConfig::default());
@@ -445,6 +514,46 @@ mod tests {
}
}
+ #[cfg(feature = "hdfs-object-store")]
+ #[test]
+ fn test_hdfs_config_serde() {
+ let config: ObjectStoreConfig = toml::from_str(
+ r#"
+type = "Hdfs"
+name = "hdfs-store"
+root = "/greptimedb"
+name_node = "hdfs://127.0.0.1:9000"
+
+[options]
+"dfs.client.block.write.replace-datanode-on-failure.enable" = "true"
+"#,
+ )
+ .unwrap();
+
+ let ObjectStoreConfig::Hdfs(hdfs_config) = config else {
+ unreachable!()
+ };
+
+ assert_eq!("hdfs-store", hdfs_config.name);
+ assert_eq!("/greptimedb", hdfs_config.connection.root);
+ assert_eq!("hdfs://127.0.0.1:9000", hdfs_config.connection.name_node);
+ assert_eq!(
+ Some(&"true".to_string()),
+ hdfs_config
+ .connection
+ .options
+ .get("dfs.client.block.write.replace-datanode-on-failure.enable")
+ );
+
+ let serialized = toml::to_string(&hdfs_config).unwrap();
+ assert!(serialized.contains("name_node = \"hdfs://127.0.0.1:9000\""));
+ assert!(
+ serialized.contains(
+ "\"dfs.client.block.write.replace-datanode-on-failure.enable\" = \"true\""
+ )
+ );
+ }
+
#[cfg(feature = "mysql-object-store")]
#[test]
fn test_mysql_config_connection_string_serde() {
diff --git a/src/object-store/src/factory.rs b/src/object-store/src/factory.rs
index 9bdd656a07c..356d3aecb21 100644
--- a/src/object-store/src/factory.rs
+++ b/src/object-store/src/factory.rs
@@ -15,15 +15,21 @@
use std::{fs, path};
use common_telemetry::info;
+#[cfg(feature = "hdfs-object-store")]
+use opendal::services::HdfsNative;
#[cfg(feature = "mysql-object-store")]
use opendal::services::Mysql;
use opendal::services::{Fs, Gcs, Oss, S3};
use snafu::prelude::*;
+#[cfg(feature = "hdfs-object-store")]
+use crate::config::HdfsConfig;
#[cfg(feature = "mysql-object-store")]
use crate::config::MysqlConfig;
use crate::config::{AzblobConfig, FileConfig, GcsConfig, ObjectStoreConfig, OssConfig, S3Config};
use crate::error::{self, Result};
+#[cfg(feature = "hdfs-object-store")]
+use crate::layers::HdfsCompatibilityLayer;
use crate::services::Azblob;
use crate::util::{build_http_context, clean_temp_dir, join_dir, normalize_dir};
use crate::{ATOMIC_WRITE_DIR, OLD_ATOMIC_WRITE_DIR, ObjectStore, util};
@@ -39,11 +45,36 @@ pub async fn new_raw_object_store(
ObjectStoreConfig::Oss(oss_config) => new_oss_object_store(oss_config).await,
ObjectStoreConfig::Azblob(azblob_config) => new_azblob_object_store(azblob_config).await,
ObjectStoreConfig::Gcs(gcs_config) => new_gcs_object_store(gcs_config).await,
+ #[cfg(feature = "hdfs-object-store")]
+ ObjectStoreConfig::Hdfs(hdfs_config) => new_hdfs_object_store(hdfs_config).await,
#[cfg(feature = "mysql-object-store")]
ObjectStoreConfig::Mysql(mysql_config) => new_mysql_object_store(mysql_config).await,
}
}
+/// Creates an object store backed by a native HDFS client.
+#[cfg(feature = "hdfs-object-store")]
+pub async fn new_hdfs_object_store(hdfs_config: &HdfsConfig) -> Result {
+ let root = util::normalize_dir(&hdfs_config.connection.root);
+ info!(
+ "The HDFS NameNode is: {}, root is: {}",
+ hdfs_config.connection.name_node, root
+ );
+
+ let builder = HdfsNative::from(&hdfs_config.connection);
+ let compatibility_layer = HdfsCompatibilityLayer::new(
+ &hdfs_config.connection.name_node,
+ &hdfs_config.connection.root,
+ &hdfs_config.connection.options,
+ )
+ .context(error::InitBackendSnafu)?;
+ let operator = ObjectStore::new(builder)
+ .context(error::InitBackendSnafu)?
+ .layer(compatibility_layer);
+
+ Ok(operator)
+}
+
#[cfg(feature = "mysql-object-store")]
pub async fn new_mysql_object_store(mysql_config: &MysqlConfig) -> Result {
let root = util::normalize_dir(&mysql_config.root);
@@ -144,3 +175,33 @@ pub async fn new_s3_object_store(s3_config: &S3Config) -> Result {
Ok(operator)
}
+
+#[cfg(all(test, feature = "hdfs-object-store"))]
+mod tests {
+ use opendal::services::HDFS_NATIVE_SCHEME;
+
+ use super::*;
+ use crate::config::HdfsConnection;
+
+ #[tokio::test]
+ async fn test_new_hdfs_object_store() {
+ let config = HdfsConfig {
+ connection: HdfsConnection {
+ root: "/greptimedb".to_string(),
+ name_node: "hdfs://127.0.0.1:9000".to_string(),
+ ..Default::default()
+ },
+ ..Default::default()
+ };
+
+ let store = new_hdfs_object_store(&config).await.unwrap();
+ assert_eq!(HDFS_NATIVE_SCHEME, store.info().scheme());
+ assert_eq!("/greptimedb/", store.info().root());
+ }
+
+ #[tokio::test]
+ async fn test_new_hdfs_object_store_requires_name_node() {
+ let result = new_hdfs_object_store(&HdfsConfig::default()).await;
+ assert!(result.is_err());
+ }
+}
diff --git a/src/object-store/src/layers.rs b/src/object-store/src/layers.rs
index cc9fa9f4df5..77c894dbfcd 100644
--- a/src/object-store/src/layers.rs
+++ b/src/object-store/src/layers.rs
@@ -12,9 +12,13 @@
// See the License for the specific language governing permissions and
// limitations under the License.
+#[cfg(feature = "hdfs-object-store")]
+mod hdfs;
#[cfg(feature = "testing")]
pub mod mock;
+#[cfg(feature = "hdfs-object-store")]
+pub use hdfs::HdfsCompatibilityLayer;
pub use opendal::layers::*;
pub use prometheus::build_prometheus_metrics_layer;
diff --git a/src/object-store/src/layers/hdfs.rs b/src/object-store/src/layers/hdfs.rs
new file mode 100644
index 00000000000..4e00098fbf7
--- /dev/null
+++ b/src/object-store/src/layers/hdfs.rs
@@ -0,0 +1,551 @@
+// Copyright 2023 Greptime Team
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+use std::fmt::{self, Debug};
+use std::sync::Arc;
+
+use hdfs_native::{Client, ClientBuilder};
+use opendal::raw::oio::{Delete as _, Read as _, ReadStream as _, Write as _};
+use opendal::raw::{
+ Layer, OpCompose, OpCopy, OpCreateDir, OpDelete, OpList, OpPresign, OpRead, OpRename, OpStat,
+ OpWrite, RpCreateDir, RpPresign, RpRename, RpStat, Service, ServiceInfo, Servicer, oio,
+};
+use opendal::{Buffer, Capability, ErrorKind, Metadata, OperationContext, Result};
+use uuid::Uuid;
+
+/// Adds atomic writes and streaming copies to the native HDFS backend.
+#[derive(Debug, Clone)]
+pub struct HdfsCompatibilityLayer {
+ renamer: AtomicRenamer,
+}
+
+impl HdfsCompatibilityLayer {
+ /// Creates a compatibility layer for the HDFS connection.
+ pub fn new(
+ name_node: &str,
+ root: &str,
+ options: &std::collections::HashMap,
+ ) -> Result {
+ let mut config = std::collections::HashMap::new();
+ let namenodes = name_node
+ .split(',')
+ .filter_map(|value| {
+ let value = value
+ .trim()
+ .trim_start_matches("hdfs://")
+ .trim_end_matches('/');
+ (!value.is_empty()).then_some(value)
+ })
+ .collect::>();
+ for (index, namenode) in namenodes.iter().enumerate() {
+ config.insert(
+ format!("dfs.namenode.rpc-address.nameservice.nn{index}"),
+ (*namenode).to_string(),
+ );
+ }
+ config.insert(
+ "dfs.ha.namenodes.nameservice".to_string(),
+ (0..namenodes.len())
+ .map(|index| format!("nn{index}"))
+ .collect::>()
+ .join(","),
+ );
+ config.extend(options.clone());
+ let client = ClientBuilder::new()
+ .with_url("hdfs://nameservice")
+ .with_config(config)
+ .build()
+ .map_err(hdfs_error)?;
+ Ok(Self {
+ renamer: AtomicRenamer::Native {
+ client,
+ root: opendal::raw::normalize_root(root),
+ },
+ })
+ }
+
+ /// Creates a compatibility layer backed by the inner service's rename.
+ #[cfg(any(test, feature = "testing"))]
+ pub fn new_for_test() -> Self {
+ Self {
+ renamer: AtomicRenamer::Raw,
+ }
+ }
+}
+
+#[derive(Debug, Clone)]
+enum AtomicRenamer {
+ Native {
+ client: Client,
+ root: String,
+ },
+ #[cfg(any(test, feature = "testing"))]
+ Raw,
+}
+
+impl AtomicRenamer {
+ async fn rename(
+ &self,
+ _inner: &Servicer,
+ _ctx: &OperationContext,
+ from: &str,
+ to: &str,
+ ) -> Result<()> {
+ match self {
+ Self::Native { client, root } => {
+ // OpenDAL's HDFS rename removes an existing destination before
+ // renaming. Use HDFS Rename2 with overwrite to keep replacement atomic.
+ client
+ .rename(
+ &opendal::raw::build_rooted_abs_path(root, from),
+ &opendal::raw::build_rooted_abs_path(root, to),
+ true,
+ )
+ .await
+ .map_err(hdfs_error)
+ }
+ #[cfg(any(test, feature = "testing"))]
+ Self::Raw => _inner
+ .rename(_ctx, from, to, OpRename::new())
+ .await
+ .map(|_| ()),
+ }
+ }
+}
+
+fn hdfs_error(error: hdfs_native::HdfsError) -> opendal::Error {
+ opendal::Error::new(ErrorKind::Unexpected, "native HDFS operation failed").set_source(error)
+}
+
+impl Layer for HdfsCompatibilityLayer {
+ fn apply_service(&self, inner: Servicer) -> Servicer {
+ Arc::new(HdfsCompatibilityService {
+ inner,
+ renamer: self.renamer.clone(),
+ })
+ }
+}
+
+struct HdfsCompatibilityService {
+ inner: Servicer,
+ renamer: AtomicRenamer,
+}
+
+impl Debug for HdfsCompatibilityService {
+ fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
+ f.debug_struct("HdfsCompatibilityService")
+ .field("inner", &self.inner)
+ .finish()
+ }
+}
+
+/// A writer that publishes non-append writes with an atomic rename.
+struct HdfsWriter(HdfsWriterInner);
+
+enum HdfsWriterInner {
+ Direct(oio::Writer),
+ Atomic {
+ inner: Servicer,
+ context: OperationContext,
+ renamer: AtomicRenamer,
+ writer: Option,
+ temporary_path: String,
+ target_path: String,
+ },
+}
+
+impl oio::Write for HdfsWriter {
+ async fn write(&mut self, buffer: Buffer) -> Result<()> {
+ match &mut self.0 {
+ HdfsWriterInner::Direct(writer) => writer.write(buffer).await,
+ HdfsWriterInner::Atomic { writer, .. } => {
+ writer
+ .as_mut()
+ .ok_or_else(writer_unavailable)?
+ .write(buffer)
+ .await
+ }
+ }
+ }
+
+ async fn close(&mut self) -> Result {
+ match &mut self.0 {
+ HdfsWriterInner::Direct(writer) => writer.close().await,
+ HdfsWriterInner::Atomic {
+ inner,
+ context,
+ renamer,
+ writer,
+ temporary_path,
+ target_path,
+ } => {
+ let mut writer = writer.take().ok_or_else(writer_unavailable)?;
+ let metadata = match writer.close().await {
+ Ok(metadata) => metadata,
+ Err(error) => {
+ drop(writer);
+ let _ = delete_path(inner, context, temporary_path).await;
+ return Err(error);
+ }
+ };
+ drop(writer);
+
+ if let Err(error) = renamer
+ .rename(inner, context, temporary_path, target_path)
+ .await
+ {
+ let _ = delete_path(inner, context, temporary_path).await;
+ return Err(error);
+ }
+
+ Ok(metadata)
+ }
+ }
+ }
+
+ async fn abort(&mut self) -> Result<()> {
+ match &mut self.0 {
+ HdfsWriterInner::Direct(writer) => writer.abort().await,
+ HdfsWriterInner::Atomic {
+ inner,
+ context,
+ writer,
+ temporary_path,
+ ..
+ } => {
+ let abort_result = if let Some(mut writer) = writer.take() {
+ let result = writer.abort().await;
+ drop(writer);
+ result
+ } else {
+ Ok(())
+ };
+ let cleanup_result = delete_path(inner, context, temporary_path).await;
+
+ match abort_result {
+ Err(error) if error.kind() != ErrorKind::Unsupported => Err(error),
+ _ => cleanup_result,
+ }
+ }
+ }
+ }
+}
+
+fn writer_unavailable() -> opendal::Error {
+ opendal::Error::new(
+ ErrorKind::Unexpected,
+ "HDFS writer is unavailable after close or abort",
+ )
+}
+
+impl Service for HdfsCompatibilityService {
+ type Reader = oio::Reader;
+ type Writer = HdfsWriter;
+ type Lister = oio::Lister;
+ type Deleter = oio::Deleter;
+ type Copier = oio::OneShotCopier;
+ type Composer = oio::Composer;
+
+ fn info(&self) -> ServiceInfo {
+ self.inner.info()
+ }
+
+ fn capability(&self) -> Capability {
+ let mut capability = self.inner.capability();
+ capability.copy = true;
+ capability
+ }
+
+ async fn create_dir(
+ &self,
+ ctx: &OperationContext,
+ path: &str,
+ args: OpCreateDir,
+ ) -> Result {
+ self.inner.create_dir(ctx, path, args).await
+ }
+
+ async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result {
+ self.inner.stat(ctx, path, args).await
+ }
+
+ fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result {
+ self.inner.read(ctx, path, args)
+ }
+
+ fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result {
+ if args.append() {
+ return self
+ .inner
+ .write(ctx, path, args)
+ .map(|writer| HdfsWriter(HdfsWriterInner::Direct(writer)));
+ }
+
+ let temporary_path = temporary_path(path);
+ let writer = self.inner.write(ctx, &temporary_path, args)?;
+ Ok(HdfsWriter(HdfsWriterInner::Atomic {
+ inner: Arc::clone(&self.inner),
+ context: ctx.clone(),
+ renamer: self.renamer.clone(),
+ writer: Some(writer),
+ temporary_path,
+ target_path: path.to_string(),
+ }))
+ }
+
+ fn copy(
+ &self,
+ ctx: &OperationContext,
+ from: &str,
+ to: &str,
+ args: OpCopy,
+ ) -> Result {
+ if args.if_not_exists() || args.if_match().is_some() {
+ return Err(opendal::Error::new(
+ ErrorKind::Unsupported,
+ "conditional copy is not supported by the HDFS fallback",
+ ));
+ }
+
+ let inner = Arc::clone(&self.inner);
+ let context = ctx.clone();
+ let renamer = self.renamer.clone();
+ let from = from.to_string();
+ let to = to.to_string();
+ Ok(oio::OneShotCopier::new(async move {
+ copy_via_read_write(inner, &context, renamer, &from, &to).await
+ }))
+ }
+
+ fn delete(&self, ctx: &OperationContext) -> Result {
+ self.inner.delete(ctx)
+ }
+
+ fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result {
+ self.inner.list(ctx, path, args)
+ }
+
+ fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result {
+ self.inner.compose(ctx, to, args)
+ }
+
+ async fn rename(
+ &self,
+ ctx: &OperationContext,
+ from: &str,
+ to: &str,
+ args: OpRename,
+ ) -> Result {
+ self.inner.rename(ctx, from, to, args).await
+ }
+
+ async fn presign(
+ &self,
+ ctx: &OperationContext,
+ path: &str,
+ args: OpPresign,
+ ) -> Result {
+ self.inner.presign(ctx, path, args).await
+ }
+}
+
+async fn copy_via_read_write(
+ inner: Servicer,
+ context: &OperationContext,
+ renamer: AtomicRenamer,
+ source_path: &str,
+ target_path: &str,
+) -> Result {
+ let reader = inner.read(context, source_path, OpRead::new())?;
+ let (_, mut reader) = reader.open(opendal::BytesRange::from(..)).await?;
+ let temporary_path = temporary_path(target_path);
+ let mut writer = inner.write(context, &temporary_path, OpWrite::new())?;
+
+ loop {
+ let buffer = match reader.read().await {
+ Ok(buffer) => buffer,
+ Err(error) => {
+ abort_and_delete(&inner, context, writer, &temporary_path).await;
+ return Err(error);
+ }
+ };
+ if buffer.is_empty() {
+ break;
+ }
+ if let Err(error) = writer.write(buffer).await {
+ abort_and_delete(&inner, context, writer, &temporary_path).await;
+ return Err(error);
+ }
+ }
+
+ let metadata = match writer.close().await {
+ Ok(metadata) => metadata,
+ Err(error) => {
+ drop(writer);
+ let _ = delete_path(&inner, context, &temporary_path).await;
+ return Err(error);
+ }
+ };
+ drop(writer);
+
+ if let Err(error) = renamer
+ .rename(&inner, context, &temporary_path, target_path)
+ .await
+ {
+ let _ = delete_path(&inner, context, &temporary_path).await;
+ return Err(error);
+ }
+
+ Ok(metadata)
+}
+
+async fn abort_and_delete(
+ inner: &Servicer,
+ context: &OperationContext,
+ mut writer: oio::Writer,
+ path: &str,
+) {
+ let _ = writer.abort().await;
+ drop(writer);
+ let _ = delete_path(inner, context, path).await;
+}
+
+async fn delete_path(inner: &Servicer, context: &OperationContext, path: &str) -> Result<()> {
+ let mut deleter = inner.delete(context)?;
+ deleter.delete(path, OpDelete::new()).await?;
+ deleter.close().await
+}
+
+// TODO(fengjiachun): Clean up temporary files left behind after a process crash.
+fn temporary_path(path: &str) -> String {
+ let suffix = format!(".greptime-{}.tmp", Uuid::new_v4());
+ match path.rsplit_once('/') {
+ Some((parent, name)) => format!("{parent}/.{name}{suffix}"),
+ None => format!(".{path}{suffix}"),
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use opendal::services::Fs;
+ use opendal::{Operator, Writer};
+ use tempfile::TempDir;
+
+ use super::*;
+
+ fn test_store() -> (TempDir, Operator) {
+ let directory = tempfile::tempdir().unwrap();
+ let store = Operator::new(Fs::default().root(directory.path().to_str().unwrap()))
+ .unwrap()
+ .layer(HdfsCompatibilityLayer::new_for_test());
+ (directory, store)
+ }
+
+ #[tokio::test]
+ async fn test_service_operations() {
+ let (_directory, store) = test_store();
+ store.create_dir("data/").await.unwrap();
+ store.write("data/source", "contents").await.unwrap();
+ assert_eq!(8, store.stat("data/source").await.unwrap().content_length());
+ assert_eq!(
+ b"onte",
+ store
+ .read_with("data/source")
+ .range(1..5)
+ .await
+ .unwrap()
+ .to_bytes()
+ .as_ref()
+ );
+ store.rename("data/source", "data/target").await.unwrap();
+ let entries = store.list("data/").await.unwrap();
+ assert!(entries.iter().any(|entry| entry.path() == "data/target"));
+ assert!(!store.exists("data/source").await.unwrap());
+ store.delete("data/target").await.unwrap();
+ assert!(!store.exists("data/target").await.unwrap());
+ }
+
+ #[tokio::test]
+ async fn test_atomic_write_keeps_old_data_after_abort() {
+ let (_directory, store) = test_store();
+ store.write("manifest.json", "old").await.unwrap();
+
+ let mut writer: Writer = store.writer("manifest.json").await.unwrap();
+ writer.write("new").await.unwrap();
+ writer.abort().await.unwrap();
+
+ assert_eq!(
+ b"old",
+ store
+ .read("manifest.json")
+ .await
+ .unwrap()
+ .to_bytes()
+ .as_ref()
+ );
+ assert!(
+ store
+ .list("")
+ .await
+ .unwrap()
+ .iter()
+ .all(|entry| !entry.path().contains(".greptime-"))
+ );
+ }
+
+ #[tokio::test]
+ async fn test_atomic_write_replaces_on_close() {
+ let (_directory, store) = test_store();
+ store.write("manifest.json", "old").await.unwrap();
+ store.write("manifest.json", "new").await.unwrap();
+
+ assert_eq!(
+ b"new",
+ store
+ .read("manifest.json")
+ .await
+ .unwrap()
+ .to_bytes()
+ .as_ref()
+ );
+ assert!(
+ store
+ .list("")
+ .await
+ .unwrap()
+ .iter()
+ .all(|entry| !entry.path().contains(".greptime-"))
+ );
+ }
+
+ #[tokio::test]
+ async fn test_copy_fallback_streams_to_target() {
+ let (_directory, store) = test_store();
+ store.write("source.parquet", "contents").await.unwrap();
+ store
+ .copy("source.parquet", "target.parquet")
+ .await
+ .unwrap();
+
+ assert_eq!(
+ b"contents",
+ store
+ .read("target.parquet")
+ .await
+ .unwrap()
+ .to_bytes()
+ .as_ref()
+ );
+ }
+}
diff --git a/src/object-store/src/lib.rs b/src/object-store/src/lib.rs
index 2de2b570082..14a2c3d1c68 100644
--- a/src/object-store/src/lib.rs
+++ b/src/object-store/src/lib.rs
@@ -30,6 +30,8 @@ pub mod secure_fs;
pub mod test_util;
pub mod util;
+#[cfg(feature = "hdfs-object-store")]
+pub use config::HdfsConnection;
pub use config::{AzblobConnection, GcsConnection, OssConnection, S3Connection};
/// The default object cache directory name.