Skip to main content

object_store/layers/
hdfs.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::fmt::{self, Debug};
16use std::sync::Arc;
17
18use hdfs_native::{Client, ClientBuilder};
19use opendal::raw::oio::{Delete as _, Read as _, ReadStream as _, Write as _};
20use opendal::raw::{
21    Layer, OpCompose, OpCopy, OpCreateDir, OpDelete, OpList, OpPresign, OpRead, OpRename, OpStat,
22    OpWrite, RpCreateDir, RpPresign, RpRename, RpStat, Service, ServiceInfo, Servicer, oio,
23};
24use opendal::{Buffer, Capability, ErrorKind, Metadata, OperationContext, Result};
25use uuid::Uuid;
26
27/// Adds atomic writes and streaming copies to the native HDFS backend.
28#[derive(Debug, Clone)]
29pub struct HdfsCompatibilityLayer {
30    renamer: AtomicRenamer,
31}
32
33impl HdfsCompatibilityLayer {
34    /// Creates a compatibility layer for the HDFS connection.
35    pub fn new(
36        name_node: &str,
37        root: &str,
38        options: &std::collections::HashMap<String, String>,
39    ) -> Result<Self> {
40        let mut config = std::collections::HashMap::new();
41        let namenodes = name_node
42            .split(',')
43            .filter_map(|value| {
44                let value = value
45                    .trim()
46                    .trim_start_matches("hdfs://")
47                    .trim_end_matches('/');
48                (!value.is_empty()).then_some(value)
49            })
50            .collect::<Vec<_>>();
51        for (index, namenode) in namenodes.iter().enumerate() {
52            config.insert(
53                format!("dfs.namenode.rpc-address.nameservice.nn{index}"),
54                (*namenode).to_string(),
55            );
56        }
57        config.insert(
58            "dfs.ha.namenodes.nameservice".to_string(),
59            (0..namenodes.len())
60                .map(|index| format!("nn{index}"))
61                .collect::<Vec<_>>()
62                .join(","),
63        );
64        config.extend(options.clone());
65        let client = ClientBuilder::new()
66            .with_url("hdfs://nameservice")
67            .with_config(config)
68            .build()
69            .map_err(hdfs_error)?;
70        Ok(Self {
71            renamer: AtomicRenamer::Native {
72                client,
73                root: opendal::raw::normalize_root(root),
74            },
75        })
76    }
77
78    /// Creates a compatibility layer backed by the inner service's rename.
79    #[cfg(any(test, feature = "testing"))]
80    pub fn new_for_test() -> Self {
81        Self {
82            renamer: AtomicRenamer::Raw,
83        }
84    }
85}
86
87#[derive(Debug, Clone)]
88enum AtomicRenamer {
89    Native {
90        client: Client,
91        root: String,
92    },
93    #[cfg(any(test, feature = "testing"))]
94    Raw,
95}
96
97impl AtomicRenamer {
98    async fn rename(
99        &self,
100        _inner: &Servicer,
101        _ctx: &OperationContext,
102        from: &str,
103        to: &str,
104    ) -> Result<()> {
105        match self {
106            Self::Native { client, root } => {
107                // OpenDAL's HDFS rename removes an existing destination before
108                // renaming. Use HDFS Rename2 with overwrite to keep replacement atomic.
109                client
110                    .rename(
111                        &opendal::raw::build_rooted_abs_path(root, from),
112                        &opendal::raw::build_rooted_abs_path(root, to),
113                        true,
114                    )
115                    .await
116                    .map_err(hdfs_error)
117            }
118            #[cfg(any(test, feature = "testing"))]
119            Self::Raw => _inner
120                .rename(_ctx, from, to, OpRename::new())
121                .await
122                .map(|_| ()),
123        }
124    }
125}
126
127fn hdfs_error(error: hdfs_native::HdfsError) -> opendal::Error {
128    opendal::Error::new(ErrorKind::Unexpected, "native HDFS operation failed").set_source(error)
129}
130
131impl Layer for HdfsCompatibilityLayer {
132    fn apply_service(&self, inner: Servicer) -> Servicer {
133        Arc::new(HdfsCompatibilityService {
134            inner,
135            renamer: self.renamer.clone(),
136        })
137    }
138}
139
140struct HdfsCompatibilityService {
141    inner: Servicer,
142    renamer: AtomicRenamer,
143}
144
145impl Debug for HdfsCompatibilityService {
146    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
147        f.debug_struct("HdfsCompatibilityService")
148            .field("inner", &self.inner)
149            .finish()
150    }
151}
152
153/// A writer that publishes non-append writes with an atomic rename.
154struct HdfsWriter(HdfsWriterInner);
155
156enum HdfsWriterInner {
157    Direct(oio::Writer),
158    Atomic {
159        inner: Servicer,
160        context: OperationContext,
161        renamer: AtomicRenamer,
162        writer: Option<oio::Writer>,
163        temporary_path: String,
164        target_path: String,
165    },
166}
167
168impl oio::Write for HdfsWriter {
169    async fn write(&mut self, buffer: Buffer) -> Result<()> {
170        match &mut self.0 {
171            HdfsWriterInner::Direct(writer) => writer.write(buffer).await,
172            HdfsWriterInner::Atomic { writer, .. } => {
173                writer
174                    .as_mut()
175                    .ok_or_else(writer_unavailable)?
176                    .write(buffer)
177                    .await
178            }
179        }
180    }
181
182    async fn close(&mut self) -> Result<Metadata> {
183        match &mut self.0 {
184            HdfsWriterInner::Direct(writer) => writer.close().await,
185            HdfsWriterInner::Atomic {
186                inner,
187                context,
188                renamer,
189                writer,
190                temporary_path,
191                target_path,
192            } => {
193                let mut writer = writer.take().ok_or_else(writer_unavailable)?;
194                let metadata = match writer.close().await {
195                    Ok(metadata) => metadata,
196                    Err(error) => {
197                        drop(writer);
198                        let _ = delete_path(inner, context, temporary_path).await;
199                        return Err(error);
200                    }
201                };
202                drop(writer);
203
204                if let Err(error) = renamer
205                    .rename(inner, context, temporary_path, target_path)
206                    .await
207                {
208                    let _ = delete_path(inner, context, temporary_path).await;
209                    return Err(error);
210                }
211
212                Ok(metadata)
213            }
214        }
215    }
216
217    async fn abort(&mut self) -> Result<()> {
218        match &mut self.0 {
219            HdfsWriterInner::Direct(writer) => writer.abort().await,
220            HdfsWriterInner::Atomic {
221                inner,
222                context,
223                writer,
224                temporary_path,
225                ..
226            } => {
227                let abort_result = if let Some(mut writer) = writer.take() {
228                    let result = writer.abort().await;
229                    drop(writer);
230                    result
231                } else {
232                    Ok(())
233                };
234                let cleanup_result = delete_path(inner, context, temporary_path).await;
235
236                match abort_result {
237                    Err(error) if error.kind() != ErrorKind::Unsupported => Err(error),
238                    _ => cleanup_result,
239                }
240            }
241        }
242    }
243}
244
245fn writer_unavailable() -> opendal::Error {
246    opendal::Error::new(
247        ErrorKind::Unexpected,
248        "HDFS writer is unavailable after close or abort",
249    )
250}
251
252impl Service for HdfsCompatibilityService {
253    type Reader = oio::Reader;
254    type Writer = HdfsWriter;
255    type Lister = oio::Lister;
256    type Deleter = oio::Deleter;
257    type Copier = oio::OneShotCopier;
258    type Composer = oio::Composer;
259
260    fn info(&self) -> ServiceInfo {
261        self.inner.info()
262    }
263
264    fn capability(&self) -> Capability {
265        let mut capability = self.inner.capability();
266        capability.copy = true;
267        capability
268    }
269
270    async fn create_dir(
271        &self,
272        ctx: &OperationContext,
273        path: &str,
274        args: OpCreateDir,
275    ) -> Result<RpCreateDir> {
276        self.inner.create_dir(ctx, path, args).await
277    }
278
279    async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
280        self.inner.stat(ctx, path, args).await
281    }
282
283    fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
284        self.inner.read(ctx, path, args)
285    }
286
287    fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
288        if args.append() {
289            return self
290                .inner
291                .write(ctx, path, args)
292                .map(|writer| HdfsWriter(HdfsWriterInner::Direct(writer)));
293        }
294
295        let temporary_path = temporary_path(path);
296        let writer = self.inner.write(ctx, &temporary_path, args)?;
297        Ok(HdfsWriter(HdfsWriterInner::Atomic {
298            inner: Arc::clone(&self.inner),
299            context: ctx.clone(),
300            renamer: self.renamer.clone(),
301            writer: Some(writer),
302            temporary_path,
303            target_path: path.to_string(),
304        }))
305    }
306
307    fn copy(
308        &self,
309        ctx: &OperationContext,
310        from: &str,
311        to: &str,
312        args: OpCopy,
313    ) -> Result<Self::Copier> {
314        if args.if_not_exists() || args.if_match().is_some() {
315            return Err(opendal::Error::new(
316                ErrorKind::Unsupported,
317                "conditional copy is not supported by the HDFS fallback",
318            ));
319        }
320
321        let inner = Arc::clone(&self.inner);
322        let context = ctx.clone();
323        let renamer = self.renamer.clone();
324        let from = from.to_string();
325        let to = to.to_string();
326        Ok(oio::OneShotCopier::new(async move {
327            copy_via_read_write(inner, &context, renamer, &from, &to).await
328        }))
329    }
330
331    fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
332        self.inner.delete(ctx)
333    }
334
335    fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
336        self.inner.list(ctx, path, args)
337    }
338
339    fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
340        self.inner.compose(ctx, to, args)
341    }
342
343    async fn rename(
344        &self,
345        ctx: &OperationContext,
346        from: &str,
347        to: &str,
348        args: OpRename,
349    ) -> Result<RpRename> {
350        self.inner.rename(ctx, from, to, args).await
351    }
352
353    async fn presign(
354        &self,
355        ctx: &OperationContext,
356        path: &str,
357        args: OpPresign,
358    ) -> Result<RpPresign> {
359        self.inner.presign(ctx, path, args).await
360    }
361}
362
363async fn copy_via_read_write(
364    inner: Servicer,
365    context: &OperationContext,
366    renamer: AtomicRenamer,
367    source_path: &str,
368    target_path: &str,
369) -> Result<Metadata> {
370    let reader = inner.read(context, source_path, OpRead::new())?;
371    let (_, mut reader) = reader.open(opendal::BytesRange::from(..)).await?;
372    let temporary_path = temporary_path(target_path);
373    let mut writer = inner.write(context, &temporary_path, OpWrite::new())?;
374
375    loop {
376        let buffer = match reader.read().await {
377            Ok(buffer) => buffer,
378            Err(error) => {
379                abort_and_delete(&inner, context, writer, &temporary_path).await;
380                return Err(error);
381            }
382        };
383        if buffer.is_empty() {
384            break;
385        }
386        if let Err(error) = writer.write(buffer).await {
387            abort_and_delete(&inner, context, writer, &temporary_path).await;
388            return Err(error);
389        }
390    }
391
392    let metadata = match writer.close().await {
393        Ok(metadata) => metadata,
394        Err(error) => {
395            drop(writer);
396            let _ = delete_path(&inner, context, &temporary_path).await;
397            return Err(error);
398        }
399    };
400    drop(writer);
401
402    if let Err(error) = renamer
403        .rename(&inner, context, &temporary_path, target_path)
404        .await
405    {
406        let _ = delete_path(&inner, context, &temporary_path).await;
407        return Err(error);
408    }
409
410    Ok(metadata)
411}
412
413async fn abort_and_delete(
414    inner: &Servicer,
415    context: &OperationContext,
416    mut writer: oio::Writer,
417    path: &str,
418) {
419    let _ = writer.abort().await;
420    drop(writer);
421    let _ = delete_path(inner, context, path).await;
422}
423
424async fn delete_path(inner: &Servicer, context: &OperationContext, path: &str) -> Result<()> {
425    let mut deleter = inner.delete(context)?;
426    deleter.delete(path, OpDelete::new()).await?;
427    deleter.close().await
428}
429
430// TODO(fengjiachun): Clean up temporary files left behind after a process crash.
431fn temporary_path(path: &str) -> String {
432    let suffix = format!(".greptime-{}.tmp", Uuid::new_v4());
433    match path.rsplit_once('/') {
434        Some((parent, name)) => format!("{parent}/.{name}{suffix}"),
435        None => format!(".{path}{suffix}"),
436    }
437}
438
439#[cfg(test)]
440mod tests {
441    use opendal::services::Fs;
442    use opendal::{Operator, Writer};
443    use tempfile::TempDir;
444
445    use super::*;
446
447    fn test_store() -> (TempDir, Operator) {
448        let directory = tempfile::tempdir().unwrap();
449        let store = Operator::new(Fs::default().root(directory.path().to_str().unwrap()))
450            .unwrap()
451            .layer(HdfsCompatibilityLayer::new_for_test());
452        (directory, store)
453    }
454
455    #[tokio::test]
456    async fn test_service_operations() {
457        let (_directory, store) = test_store();
458        store.create_dir("data/").await.unwrap();
459        store.write("data/source", "contents").await.unwrap();
460        assert_eq!(8, store.stat("data/source").await.unwrap().content_length());
461        assert_eq!(
462            b"onte",
463            store
464                .read_with("data/source")
465                .range(1..5)
466                .await
467                .unwrap()
468                .to_bytes()
469                .as_ref()
470        );
471        store.rename("data/source", "data/target").await.unwrap();
472        let entries = store.list("data/").await.unwrap();
473        assert!(entries.iter().any(|entry| entry.path() == "data/target"));
474        assert!(!store.exists("data/source").await.unwrap());
475        store.delete("data/target").await.unwrap();
476        assert!(!store.exists("data/target").await.unwrap());
477    }
478
479    #[tokio::test]
480    async fn test_atomic_write_keeps_old_data_after_abort() {
481        let (_directory, store) = test_store();
482        store.write("manifest.json", "old").await.unwrap();
483
484        let mut writer: Writer = store.writer("manifest.json").await.unwrap();
485        writer.write("new").await.unwrap();
486        writer.abort().await.unwrap();
487
488        assert_eq!(
489            b"old",
490            store
491                .read("manifest.json")
492                .await
493                .unwrap()
494                .to_bytes()
495                .as_ref()
496        );
497        assert!(
498            store
499                .list("")
500                .await
501                .unwrap()
502                .iter()
503                .all(|entry| !entry.path().contains(".greptime-"))
504        );
505    }
506
507    #[tokio::test]
508    async fn test_atomic_write_replaces_on_close() {
509        let (_directory, store) = test_store();
510        store.write("manifest.json", "old").await.unwrap();
511        store.write("manifest.json", "new").await.unwrap();
512
513        assert_eq!(
514            b"new",
515            store
516                .read("manifest.json")
517                .await
518                .unwrap()
519                .to_bytes()
520                .as_ref()
521        );
522        assert!(
523            store
524                .list("")
525                .await
526                .unwrap()
527                .iter()
528                .all(|entry| !entry.path().contains(".greptime-"))
529        );
530    }
531
532    #[tokio::test]
533    async fn test_copy_fallback_streams_to_target() {
534        let (_directory, store) = test_store();
535        store.write("source.parquet", "contents").await.unwrap();
536        store
537            .copy("source.parquet", "target.parquet")
538            .await
539            .unwrap();
540
541        assert_eq!(
542            b"contents",
543            store
544                .read("target.parquet")
545                .await
546                .unwrap()
547                .to_bytes()
548                .as_ref()
549        );
550    }
551}