hdfs-native 0.14.2

Native HDFS client implementation in Rust
Documentation
#[cfg(feature = "integration-test")]
mod test {

    use std::{
        collections::{HashMap, HashSet},
        sync::atomic::Ordering,
        time::Duration,
    };

    use bytes::{Buf, BufMut, Bytes, BytesMut};
    use hdfs_native::{
        Client, ClientBuilder, Result, WriteOptions,
        file::FileReader,
        minidfs::{DfsFeatures, MiniDfs},
        test::{WRITE_CONNECTION_FAULT_INJECTOR, WRITE_REPLY_FAULT_INJECTOR},
    };
    use serial_test::serial;

    #[tokio::test]
    #[serial]
    async fn test_lease_renewal() -> Result<()> {
        let _ = env_logger::builder().is_test(true).try_init();

        let _dfs = MiniDfs::with_features(&HashSet::from([DfsFeatures::HA]));
        let client = Client::default();

        // Second client for checking lease
        let client2 = Client::default();

        let file = "/testfile";

        let mut writer = client.create(file, WriteOptions::default()).await?;

        writer.write_bytes(Bytes::from(vec![0u8, 1, 2, 3])).await?;

        // First client owns the lease, so can't append to the file
        assert!(client2.append("/testfile").await.is_err());

        tokio::time::sleep(Duration::from_secs(70)).await;

        // First client should still own the lease from the lease renewal. If not this would
        // trigger lease recovery on the namenode and the first client wouldn't be able to finish
        // writing.
        assert!(client2.append("/testfile").await.is_err());

        writer.write_bytes(Bytes::from(vec![4u8, 5, 6, 7])).await?;

        // This succeeding means the lease was renewed
        writer.close().await?;

        Ok(())
    }

    #[tokio::test]
    #[serial]
    async fn test_write_failures_basic() -> Result<()> {
        let _ = env_logger::builder().is_test(true).try_init();

        let _dfs = MiniDfs::with_features(&HashSet::from([DfsFeatures::HA]));
        test_write_failures().await?;
        Ok(())
    }

    #[tokio::test]
    #[serial]
    async fn test_write_failures_security() -> Result<()> {
        let _ = env_logger::builder().is_test(true).try_init();

        let _dfs = MiniDfs::with_features(&HashSet::from([DfsFeatures::HA, DfsFeatures::Security]));
        test_write_failures().await?;
        Ok(())
    }

    async fn test_write_failures() -> Result<()> {
        fn replace_dn_conf(should_replace: bool) -> HashMap<String, String> {
            HashMap::from([
                (
                    "dfs.client.block.write.replace-datanode-on-failure.enable".to_string(),
                    should_replace.to_string(),
                ),
                (
                    "dfs.client.block.write.replace-datanode-on-failure.policy".to_string(),
                    "ALWAYS".to_string(),
                ),
            ])
        }

        for (i, client) in [
            ClientBuilder::new()
                .with_config(replace_dn_conf(true))
                .build()
                .unwrap(),
            ClientBuilder::new()
                .with_config(replace_dn_conf(false))
                .build()
                .unwrap(),
        ]
        .iter()
        .enumerate()
        {
            let file = format!("/testfile{i}");
            let bytes_to_write = 2usize * 1024 * 1024;

            let mut data = BytesMut::with_capacity(bytes_to_write);
            for i in 0..(bytes_to_write / 4) {
                data.put_i32(i as i32);
            }

            // Test connection failure before writing data
            let mut writer = client
                .create(&file, WriteOptions::default().replication(3))
                .await?;

            WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);

            let data = data.freeze();
            writer.write_bytes(data.clone()).await?;
            writer.close().await?;

            let reader = client.read(&file).await?;
            check_file_content(&reader, data.clone()).await?;

            // Test connection failure after data has been written
            let mut writer = client
                .create(
                    &file,
                    WriteOptions::default().replication(3).overwrite(true),
                )
                .await?;

            writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;

            // Give a little time for the packets to send
            tokio::time::sleep(Duration::from_millis(100)).await;

            WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);

            writer.write_bytes(data.slice(bytes_to_write / 2..)).await?;
            writer.close().await?;

            let reader = client.read(&file).await?;
            check_file_content(&reader, data.clone()).await?;

            // Test failure in from ack status before any data is written
            let mut writer = client
                .create(
                    &file,
                    WriteOptions::default().replication(3).overwrite(true),
                )
                .await?;

            *WRITE_REPLY_FAULT_INJECTOR.lock().unwrap() = Some(2);

            writer.write_bytes(data.clone()).await?;
            writer.close().await?;

            let reader = client.read(&file).await?;
            check_file_content(&reader, data.clone()).await?;

            // Test failure in from ack status after some data has been written
            let mut writer = client
                .create(
                    &file,
                    WriteOptions::default().replication(3).overwrite(true),
                )
                .await?;

            writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;

            // Give a little time for the packets to send
            tokio::time::sleep(Duration::from_millis(100)).await;

            *WRITE_REPLY_FAULT_INJECTOR.lock().unwrap() = Some(2);

            writer.write_bytes(data.slice(bytes_to_write / 2..)).await?;
            writer.close().await?;

            let reader = client.read(&file).await?;
            check_file_content(&reader, data.clone()).await?;
        }

        Ok(())
    }

    #[tokio::test]
    #[serial]
    async fn test_replace_failed_datanode() -> Result<()> {
        let _ = env_logger::builder().is_test(true).try_init();

        let mut replace_dn_on_failure_conf: HashMap<String, String> = HashMap::new();
        replace_dn_on_failure_conf.insert(
            "dfs.client.block.write.replace-datanode-on-failure.enable".to_string(),
            "true".to_string(),
        );
        replace_dn_on_failure_conf.insert(
            "dfs.client.block.write.replace-datanode-on-failure.policy".to_string(),
            "ALWAYS".to_string(),
        );

        let _dfs = MiniDfs::with_features(&HashSet::from([DfsFeatures::HA]));
        let client = ClientBuilder::new()
            .with_config(replace_dn_on_failure_conf)
            .build()
            .unwrap();

        let file = "/testfile_replace_failed_datanode";
        let bytes_to_write = 2usize * 1024 * 1024;

        let mut data = BytesMut::with_capacity(bytes_to_write);
        for i in 0..(bytes_to_write / 4) {
            data.put_i32(i as i32);
        }
        let data = data.freeze();

        let mut writer = client
            .create(file, WriteOptions::default().replication(3))
            .await?;

        writer.write_bytes(data.slice(..bytes_to_write / 3)).await?;

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);
        writer
            .write_bytes(data.slice(bytes_to_write / 3..2 * bytes_to_write / 3))
            .await?;
        // Give a little time for the packets to send
        tokio::time::sleep(Duration::from_millis(100)).await;

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);
        writer
            .write_bytes(data.slice(2 * bytes_to_write / 3..))
            .await?;
        // Give a little time for the packets to send
        tokio::time::sleep(Duration::from_millis(100)).await;

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);
        writer.close().await?;
        // Give a little time for the packets to send
        tokio::time::sleep(Duration::from_millis(100)).await;

        let reader = client.read(file).await?;
        check_file_content(&reader, data).await?;

        Ok(())
    }

    async fn check_file_content(reader: &FileReader, mut expected: Bytes) -> Result<()> {
        assert_eq!(reader.file_length(), expected.len());

        let mut file_data = reader.read_range(0, reader.file_length()).await?;
        for _ in 0..expected.len() / 4 {
            assert_eq!(file_data.get_i32(), expected.get_i32());
        }
        Ok(())
    }

    #[tokio::test]
    #[serial]
    async fn test_striped_write_failures() -> Result<()> {
        let _ = env_logger::builder().is_test(true).try_init();

        let dfs_features = HashSet::from([DfsFeatures::HA, DfsFeatures::EC]);
        let _dfs = MiniDfs::with_features(&dfs_features);
        let client = Client::default();

        // A full stripe is 3 MiB, so write two full stripes to test two failures
        let bytes_to_write = 6usize * 1024 * 1024;

        let mut data = BytesMut::with_capacity(bytes_to_write);
        for i in 0..(bytes_to_write / 4) {
            data.put_i32(i as i32);
        }
        let data = data.freeze();

        // Test connection failure before writing data with erasure coding
        let file = "/ec-3-2/striped_testfile1";
        let mut writer = client.create(file, WriteOptions::default()).await?;

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);

        writer.write_bytes(data.clone()).await?;
        writer.close().await?;

        let reader = client.read(file).await?;
        check_file_content(&reader, data.clone()).await?;

        // Test two connection failures after data has been written
        let file = "/ec-3-2/striped_testfile2";
        let mut writer = client.create(file, WriteOptions::default()).await?;

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);

        writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;

        // Give a little time for the packets to send
        tokio::time::sleep(Duration::from_millis(100)).await;

        assert!(!WRITE_CONNECTION_FAULT_INJECTOR.load(Ordering::SeqCst));
        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);

        writer.write_bytes(data.slice(bytes_to_write / 2..)).await?;
        writer.close().await?;

        // Give a little time for the packets to send
        tokio::time::sleep(Duration::from_millis(100)).await;

        assert!(!WRITE_CONNECTION_FAULT_INJECTOR.load(Ordering::SeqCst));

        let reader = client.read(file).await?;
        check_file_content(&reader, data.clone()).await?;

        // Test failure in from ack status
        let file = "/ec-3-2/striped_testfile3";
        let mut writer = client.create(file, WriteOptions::default()).await?;

        *WRITE_REPLY_FAULT_INJECTOR.lock().unwrap() = Some(2);

        writer.write_bytes(data.clone()).await?;
        writer.close().await?;

        let reader = client.read(file).await?;
        check_file_content(&reader, data.clone()).await?;

        // Test failure in from ack status after some data has been written
        let file = "/ec-3-2/striped_testfile4";
        let mut writer = client.create(file, WriteOptions::default()).await?;

        writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;

        // Give a little time for the packets to send
        tokio::time::sleep(Duration::from_millis(100)).await;

        *WRITE_REPLY_FAULT_INJECTOR.lock().unwrap() = Some(2);

        writer.write_bytes(data.slice(bytes_to_write / 2..)).await?;
        writer.close().await?;

        let reader = client.read(file).await?;
        check_file_content(&reader, data.clone()).await?;

        // Test three failures which should cause a write failure
        let file = "/ec-3-2/striped_testfile5";
        let mut writer = client.create(file, WriteOptions::default()).await?;

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);
        writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;
        tokio::time::sleep(Duration::from_millis(100)).await;
        assert!(!WRITE_CONNECTION_FAULT_INJECTOR.load(Ordering::SeqCst));

        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);
        writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;
        tokio::time::sleep(Duration::from_millis(100)).await;
        assert!(!WRITE_CONNECTION_FAULT_INJECTOR.load(Ordering::SeqCst));

        // Third write. It won't fail yet because the sending of the data is asynchronous.
        // The failure will happen when the client tries to send the next block.
        WRITE_CONNECTION_FAULT_INJECTOR.store(true, Ordering::SeqCst);
        writer.write_bytes(data.slice(..bytes_to_write / 2)).await?;
        tokio::time::sleep(Duration::from_millis(100)).await;
        assert!(!WRITE_CONNECTION_FAULT_INJECTOR.load(Ordering::SeqCst));

        assert!(
            writer
                .write_bytes(data.slice(..bytes_to_write / 2))
                .await
                .is_err()
        );

        Ok(())
    }
}