All Products
Search
Document Center

E-MapReduce:Separate hot and cold data with HDFS

Last Updated:Sep 21, 2026

This topic shows how to implement hot and cold data separation on a ClickHouse cluster in Alibaba Cloud E-MapReduce using HDFS. This method helps you automatically manage data tiers, optimizing storage costs while maintaining cluster read and write performance.

Prerequisites

  • You have created a ClickHouse cluster of EMR V5.5.0 or later. For more information, see Create a ClickHouse cluster.

  • An HDFS service, such as one from an EMR Hadoop cluster, is available in the same Virtual Private Cloud (VPC).

  • You have read and write permissions for the HDFS service.

Limitations

This procedure applies only to ClickHouse clusters of EMR V5.5.0 or later.

Procedure

  1. Step 1: Add an HDFS disk in the EMR console

  2. Step 2: Verify the configuration

  3. Step 3: Implement hot and cold data separation

Step 1: Add an HDFS disk

  1. Go to the configuration page of the ClickHouse service.

    1. Log on to the EMR on ECS console.

    2. In the top navigation bar, select the region and a resource group.

    3. On the EMR on ECS page, find your cluster and click Services in the Actions column.

    4. On the Services page, click Configure in the ClickHouse section.

  2. On the Configure tab, click the server-metrika tab.

  3. Modify the value of the storage_configuration parameter.

    1. In disks, add an HDFS disk.

      The following provides a sample configuration.

      <disk_hdfs>
          <type>hdfs</type>
          <endpoint>hdfs://${your-hdfs-url}</endpoint>
          <min_bytes_for_seek>1048576</min_bytes_for_seek>
          <thread_pool_size>16</thread_pool_size>
          <objects_chunk_size_to_delete>1000</objects_chunk_size_to_delete>
      </disk_hdfs>

      The following table describes the parameters.

      Parameter

      Required

      Description

      disk_hdfs

      Yes

      The name of the disk. You can specify a custom name.

      type

      Yes

      The disk type. Set this to hdfs.

      endpoint

      Yes

      The endpoint of the HDFS service.

      Important

      The HDFS endpoint is typically the address of the NameNode. If the NameNode is in High Availability (HA) mode, the port is usually 8020. Otherwise, the port is 9000.

      min_bytes_for_seek

      No

      The minimum number of bytes for a seek operation. If a seek offset is smaller than this value, the skip operation is used instead. The default value is 1048576.

      thread_pool_size

      No

      The size of the thread pool used by the disk to perform restore operations. The default value is 16.

      objects_chunk_size_to_delete

      No

      The maximum number of HDFS files that can be deleted in a single operation. The default value is 1000.

    2. In policies, add a new policy.

      The following provides a sample policy.

      <hdfs_ttl>
          <volumes>
            <local>
              <!-- Include all disks under the default storage policy. -->
              <disk>disk1</disk>
              <disk>disk2</disk>
              <disk>disk3</disk>
              <disk>disk4</disk>
            </local>
            <remote>
              <disk>disk_hdfs</disk>
            </remote>
          </volumes>
          <move_factor>0.2</move_factor>
      </hdfs_ttl>
      Note

      You can also add this configuration directly to the default policy.

  4. Save the configuration.

    1. On the Configure page of the ClickHouse service, click Save.

    2. In the Modify Information dialog box, enter a reason, enable the Save and Deliver Configuration switch, and click Save.

  5. Deploy the client configuration.

    1. On the Configure page of the ClickHouse service, click Deploy Client Configuration.

    2. In the Deploy Client Configuration dialog box, enter a reason and click OK.

    3. In the Confirm dialog box, click OK.

Step 2: Verify the configuration

  1. Log on to the ClickHouse cluster using SSH. For more information, see Log on to a cluster.

  2. Run the following command to start the ClickHouse client:

    clickhouse-client -h core-1-1 -m
    Note

    This example logs on to the core-1-1 node. If you have multiple core nodes, you can log on to any one of them.

  3. Run the following command to view disk information.

    select * from system.disks;

    The following output is returned.

    ┌─name─────┬─path─────────────────────────────────┬───────────free_space─┬──────────total_space─┬─keep_free_space─┬─type──┐
    │ default  │ /var/lib/clickhouse/                 │          83868921856 │          84014424064 │               0 │ local │
    │ disk1    │ /mnt/disk1/clickhouse/               │          83858436096 │          84003938304 │        10485760 │ local │
    │ disk2    │ /mnt/disk2/clickhouse/               │          83928215552 │          84003938304 │        10485760 │ local │
    │ disk3    │ /mnt/disk3/clickhouse/               │          83928301568 │          84003938304 │        10485760 │ local │
    │ disk4    │ /mnt/disk4/clickhouse/               │          83928301568 │          84003938304 │        10485760 │ local │
    │ disk_hdfs│ /var/lib/clickhouse/disks/disk_hdfs/ │ 18446744073709551615 │ 18446744073709551615 │               0 │ hdfs  │
    └──────────┴──────────────────────────────────────┴──────────────────────┴──────────────────────┴─────────────────┴───────┘
                                
  4. Run the following command to view disk storage policies.

    select * from system.storage_policies;

    The following output is returned.

    ┌─policy_name──┬─volume_name─┬─volume_priority─┬─disks──────────────────────────────┬─volume_type─┬─max_data_part_size─┬─move_factor─┬─prefer_not_to_merge─┐
    │ default      │ single      │               1 │ ['disk1','disk2','disk3','disk4']           │JBOD        │                  0 │           0 │                   0 │
    │ hdfs_ttl     │ local       │               1 │ ['disk1','disk2','disk3','disk4']           │JBOD        │                  0 │          0.2 │                   0 │
    │ hdfs_ttl     │ remote      │               2 │ ['disk_hdfs']                         │JBOD        │                  0 │          0.2 │                   0 │
    └──────────────┴─────────────┴─────────────────┴───────────────────────────────────┴─────────────┴────────────────────┴─────────────┴─────────────────────┘

    This output indicates that the disk configuration is successfully updated.

Step 3: Separate hot and cold data

Modify an existing table

  1. View the current storage policy.

    1. In the ClickHouse client, run the following command to view disk information.

      SELECT
        storage_policy
      FROM system.tables
      WHERE database='<database_name>' AND name='<table_name>';

      In the command, <database_name> is the database name, and <table_name> is the table name.

      If the following output is returned, proceed to the next step to add a volume.

      <default>
        <volumes>
          <single>
            <disk>disk1</disk>
            <disk>disk2</disk>
            <disk>disk3</disk>
            <disk>disk4</disk>
          </single>
        </volumes>
      </default>
  2. Extend the current storage policy.

    On the Configure tab of the ClickHouse service in the EMR console, add the remote volume configuration. The following is an example:

    <default>
      <volumes>
        <single>
          <disk>disk1</disk>
          <disk>disk2</disk>
          <disk>disk3</disk>
          <disk>disk4</disk>
        </single>
        <!-- The following is the newly added remote volume. -->
        <remote>
          <disk>disk_hdfs</disk>
        </remote>
      </volumes>
      <!-- You must specify move_factor when using multiple volumes. -->
      <move_factor>0.2</move_factor>
    </default>
  3. Run the following command to modify the TTL.

    ALTER TABLE <yourDataName>.<yourTableName>
      MODIFY TTL toStartOfMinute(addMinutes(t, 5)) TO VOLUME 'remote';
  4. Run the following command to view the data part distribution.

    select partition,name,path from system.parts where database='<yourDataName>' and table='<yourTableName>' and active=1

    The following output is returned.

    ┌─partition───────────┬─name─────────────────┬─path──────────────────────────────────────────────────────────────────────────────────────────────────┐
    │ 2022-01-12 11:30:00 │ 1641958200_1_96_3    │ /var/lib/clickhouse/disks/disk_hdfs/store/156/156008ff-41bf-460c-8848-e34fad88c25d/1641958200_1_96_3/ │
    │ 2022-01-12 11:35:00 │ 1641958500_97_124_2  │ /mnt/disk3/clickhouse/store/156/156008ff-41bf-460c-8848-e34fad88c25d/1641958500_97_124_2/             │
    │ 2022-01-12 11:35:00 │ 1641958500_125_152_2 │ /mnt/disk4/clickhouse/store/156/156008ff-41bf-460c-8848-e34fad88c25d/1641958500_125_152_2/            │
    │ 2022-01-12 11:35:00 │ 1641958500_153_180_2 │ /mnt/disk1/clickhouse/store/156/156008ff-41bf-460c-8848-e34fad88c25d/1641958500_153_180_2/            │
    │ 2022-01-12 11:35:00 │ 1641958500_181_186_1 │ /mnt/disk4/clickhouse/store/156/156008ff-41bf-460c-8848-e34fad88c25d/1641958500_181_186_1/            │
    │ 2022-01-12 11:35:00 │ 1641958500_187_192_1 │ /mnt/disk3/clickhouse/store/156/156008ff-41bf-460c-8848-e34fad88c25d/1641958500_187_192_1/            │
    └─────────────────────┴──────────────────────┴───────────────────────────────────────────────────────────────────────────────────────────────────────┘
    
    6 rows in set. Elapsed: 0.002 sec.
    Note

    This output shows that hot and cold data is separated based on the TTL. Hot data is stored on local disks, and cold data is stored in HDFS.

    In the output, /var/lib/clickhouse/disks/disk_hdfs is the metadata directory for disk_hdfs, and /mnt/disk{1..4}/clickhouse is the path for the local disks.

Create a new table

  • Syntax

    CREATE TABLE <yourDataName>.<yourTableName> [ON CLUSTER cluster_emr]
    (
      column1 Type1,
      column2 Type2,
      ...
    ) Engine = MergeTree() -- You can also use a Replicated*MergeTree engine.
    PARTITION BY <yourPartitionKey>
    ORDER BY <yourPartitionKey>
    TTL <yourTtlKey> TO VOLUME 'remote'
    SETTINGS storage_policy='hdfs_ttl';
    Note

    In the preceding command, <yourPartitionKey> indicates the partition key for a table in the ClickHouse cluster. <yourTtlKey> indicates the TTL that you specify.

  • Example

    CREATE TABLE test.test
    (
        `id` UInt32,
        `t` DateTime
    )
    ENGINE = MergeTree()
    PARTITION BY toStartOfFiveMinute(t)
    ORDER BY id
    TTL toStartOfMinute(addMinutes(t, 5)) TO VOLUME 'remote'
    SETTINGS storage_policy='hdfs_ttl';
    Note

    In this example, the table stores data from the last five minutes on local disks. After five minutes, the data is moved to the remote volume, which is HDFS.

Related configurations

  • server-config

    merge_tree.allow_remote_fs_zero_copy_replication: Set this parameter to true to enable zero-copy replication for Replicated*MergeTree tables that use remote storage such as DiskHDFS. When enabled, ClickHouse uses the native replication of the remote file system. Instead of copying data files, only the metadata is replicated across different replicas within the same shard.

  • server-users

    profile.${your-profile-name}.hdfs_replication: The number of data replicas to store in HDFS.