From 75bd8e9ce6fa554c2da3fa7d72dc63c8b50e99ce Mon Sep 17 00:00:00 2001 From: Minghan Jiang Date: Mon, 28 Sep 2026 08:39:48 +0000 Subject: [PATCH] feat: add HDFS object storage backend (#8701) * feat: add HDFS object storage backend Signed-off-by: Minghan2005 * fix: make HDFS storage operations durable Gate the native HDFS backend behind an explicit feature. Publish writes through same-directory temporary files and atomic HDFS Rename2 replacement, and provide streaming copy fallback for COPY_REGION. Add regression coverage for interrupted writes and the region-copy path. Signed-off-by: Minghan2005 * ci: run HDFS object store tests Signed-off-by: jeremyhi * docs: note HDFS temporary file cleanup follow-up Signed-off-by: jeremyhi * feat: enable HDFS object storage by default Signed-off-by: jeremyhi * docs: remove redundant HDFS build feature notes Signed-off-by: jeremyhi --------- Signed-off-by: Minghan2005 Signed-off-by: jeremyhi Co-authored-by: Minghan2005 Co-authored-by: jeremyhi --- .github/workflows/nightly-ci.yml | 2 +- .github/workflows/rust.yml | 4 +- Cargo.lock | 241 ++++++++++-- Cargo.toml | 1 + config/config.md | 12 +- config/datanode.example.toml | 22 +- config/standalone.example.toml | 22 +- src/cmd/Cargo.toml | 2 + src/cmd/tests/load_config_test.rs | 38 ++ src/mito2/Cargo.toml | 1 + src/mito2/src/region/utils.rs | 63 ++++ src/object-store/Cargo.toml | 2 + src/object-store/src/config.rs | 109 ++++++ src/object-store/src/factory.rs | 61 +++ src/object-store/src/layers.rs | 4 + src/object-store/src/layers/hdfs.rs | 551 ++++++++++++++++++++++++++++ src/object-store/src/lib.rs | 2 + 17 files changed, 1090 insertions(+), 47 deletions(-) create mode 100644 src/object-store/src/layers/hdfs.rs 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.