This topic describes how to use JindoDistCp.
What is JindoDistCp
JindoDistCp is a distributed file copy tool developed by the Alibaba Cloud Data Lake Storage team for large-scale data transfers within and between clusters. It uses MapReduce to distribute files, handle errors, and recover from failures. It takes a list of files and directories as input for MapReduce tasks, where each task copies a portion of the source list. It fully supports data copy scenarios between the Hadoop Distributed File System (HDFS), OSS-HDFS, OSS, and S3. It provides various custom copy parameters and strategies. It is optimized for copying data from HDFS to OSS-HDFS. Using a custom CopyCommitter, it performs No-Rename copies and ensures data consistency upon completion. Its features are fully aligned with those of S3 DistCp and HDFS DistCp. It offers significant performance improvements over HDFS DistCp. JindoDistCp is designed to be an efficient, stable, and secure data copy tool.
Environment requirements
JDK 1.8.0 or later.
Hadoop 2.3 or later. You must download the latest version of the `jindo-distcp-tool-x.x.x.jar` file. This JAR file is included in the `jindosdk-${version}.tar.gz` package. After you decompress the package, you can find the JAR file in the `tools/` directory. For more information, see JindoData downloads.
NoteJindoDistCp is deployed on clusters that run EMR V5.6.0 or later and EMR V3.40.0 or later. You can find the `jindo-distcp-tool-x.x.x.jar` file in the `/opt/apps/JINDOSDK/jindosdk-current/tools` directory.
Parameters
JindoDistCp is provided as a JAR package. You can use the `hadoop jar` command with a series of parameters to perform migration operations.
Parameter | Parameter type | Description | Default value | Version | OSS | OSS-HDFS |
Required | Specifies the source directory. The following prefixes are supported:
| None | 4.3.0 or later | Support | Supported | |
Required | Specifies the destination directory. The following prefixes are supported:
| None | 4.3.0 or later | Supported | Support | |
Optional | Specifies the bandwidth limit for a single node. Unit: MB. | -1 | 4.3.0 or later | Support | Support | |
Optional | Specifies the compression type. Supported codecs include gzip, gz, lzo, lzop, and snappy. | keep (The compression type is not changed.) | 4.3.0 or later | Support | Supported | |
Optional | Specifies the storage policy for the destination. Valid values: Standard, IA, Archive, and ColdArchive. | Standard | 4.3.0 or later | Support | Not supported | |
Optional | Specifies the file that contains filter rules. | None | 4.3.0 or later | Support | Support | |
Optional | The settings apply to files that match the rule. | None | 4.3.0 or later | Supported | Support | |
Optional | Specifies the concurrency of the DistCp task. This corresponds to the mapreduce.job.maps parameter in a MapReduce task. | 10 | 4.3.0 or later | Support | Support | |
Optional | Specifies the number of files to be processed by each DistCp job. | 10000 | 4.5.1 or later | Support | Supported | |
Optional | Specifies the number of files to be processed by each DistCp task. | 1 | 4.3.0 or later | Supported | Support | |
Optional | Specifies the temporary directory. | /tmp | 4.3.0 or later | Supported | Support | |
Optional | Sets a configuration. | None | 4.3.0 or later | Support | Support | |
Optional | Specifies whether to disable checksum verification. | false | 4.3.0 or later | Supported | Support | |
Optional | Specifies whether to delete the source files. This is used to move data. | false | 4.3.0 or later | Supported | Supported | |
Optional | Specifies whether to enable transactions to ensure job-level atomicity. | false | 4.3.0 or later | Supported | Supported | |
Optional | Specifies whether to ignore exceptions thrown during the copy task to prevent the task from being interrupted. | false | 4.3.0 or later | Support | Support | |
Optional | Specifies whether to enable monitoring and alerting. | false | 4.5.1 or later | Support | Support | |
Optional | Sets the DistCp mode to DIFF to view the differences between the files in the source and destination. | DistCpMode.COPY | 4.3.0 or later | Supported | Support | |
Optional | Sets the DistCp mode to UPDATE to enable incremental synchronization. This skips identical files and directories and synchronizes only new or changed files and directories from the source to the destination. | DistCpMode.COPY | 4.3.0 or later | Supported | Support | |
Optional | Specifies whether to preserve metadata information. | false | 4.4.0 or later | Not supported | Supported |
--src and --dest (required)
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Supported |
`--src`: Specifies the path of the source files.
`--dest`: Specifies the path of the destination files.
The following command provides an example:
hadoop jar jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_tableYou can specify the destination directory using the `dest` path. For example, the preceding command copies files from `/data/hourly_table` to the `hourly_table` directory in the `example-oss-bucket` bucket. This behavior is different from that of Hadoop DistCp. By default, JindoDistCp copies all files from the source directory to the specified destination path but does not include the root directory of the source. You can specify a root directory in the destination path. If the directory does not exist, it is automatically created.
To copy a single file, you must specify a directory as the destination.
hadoop jar jindo-distcp-tool-${version}.jar --src /test.txt --dest oss://example-oss-bucket/tmpUse --bandWidth
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Support |
`--bandWidth`: Specifies the bandwidth that a single node can use for the DistCp task, in MB. This parameter prevents a single node from consuming too much bandwidth.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --bandWidth 6Use --codec
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
Source files are often stored in OSS or OSS-HDFS as uncompressed text, which is not ideal for storage costs or data analysis. You can use the `--codec` option to efficiently store data by compressing files online.
`--codec` specifies the file compression codec. It supports the gzip, gz, lzo, lzop, and snappy encoders, and the keywords `none` and `keep` (default). The keywords are described as follows:
`none`: Saves the files as uncompressed. If a source file is already compressed, JindoDistCp decompresses it.
`keep` (default): Copies the files as they are without changing their compression state.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --codec gzAfter the command is run, the files in the destination folder are compressed using the gz codec.
[root@emr-header-1 opt]# hdfs dfs -ls oss://example-oss-bucket/hourly_table/2017-02-01/03
Found 6 items
-rw-rw-rw- 1 938 2020-04-17 20:58 oss://example-oss-bucket/hourly_table/2017-02-01/03/000151.sst.gz
-rw-rw-rw- 1 1956 2020-04-17 20:58 oss://example-oss-bucket/hourly_table/2017-02-01/03/1.log.gz
-rw-rw-rw- 1 1956 2020-04-17 20:58 oss://example-oss-bucket/hourly_table/2017-02-01/03/2.log.gz
-rw-rw-rw- 1 1956 2020-04-17 20:58 oss://example-oss-bucket/hourly_table/2017-02-01/03/OPTIONS-000109.gz
-rw-rw-rw- 1 506 2020-04-17 20:58 oss://example-oss-bucket/hourly_table/2017-02-01/03/emp01.txt.gz
-rw-rw-rw- 1 506 2020-04-17 20:58 oss://example-oss-bucket/hourly_table/2017-02-01/03/emp06.txt.gzTo use the lzo codec in an open source Hadoop cluster, you must install the gplcompression native library and the hadoop-lzo package. If you do not have the required environment, you must use another compression method.
Use --filters
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Support |
`--filters`: Specifies a file that contains filter rules.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --filters filter.txtFor example, if the `filter.txt` file contains .*test.*, files whose paths contain the string "test" are not copied to OSS.
Use --srcPrefixesFile
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
`--srcPrefixesFile`: Specifies a file that contains inclusion rules.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --srcPrefixesFile prefixes.txtFor example, if the `prefixes.txt` file contains .*test.*, only files whose paths contain the string "test" are copied to OSS.
Use --parallelism
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
`--parallelism`: Specifies the `mapreduce.job.maps` parameter for the MapReduce task. The default value of this parameter in an EMR environment is 10. You can customize the value of this parameter based on your cluster resources to control the concurrency of the DistCp task.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /opt/tmp --dest oss://example-oss-bucket/tmp --parallelism 20Use --taskBatch
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Supported |
`--taskBatch`: Specifies the number of files to be processed by each DistCp task. The default value is 1.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --taskBatch 1Use --tmp
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Supported |
`--tmp`: Specifies a temporary directory in HDFS to store temporary data. The default value is `/tmp`, which corresponds to `hdfs:///tmp/`.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --tmp /tmpConfigure an AccessKey to access OSS or OSS-HDFS
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
`--hadoopConf`: If you are not in an EMR environment or if the passwordless access service has issues, you can use this option to specify an AccessKey to access OSS or OSS-HDFS.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --hadoopConf fs.oss.accessKeyId=yourkey --hadoopConf fs.oss.accessKeySecret=yoursecretTo avoid entering the AccessKey each time, you can pre-configure the AccessKey ID and AccessKey secret for OSS or OSS-HDFS in the `core-site.xml` file of Hadoop. In the EMR console, add the following configuration on the `core-site.xml` page of the Hadoop-Common service.
<configuration>
<property>
<name>fs.oss.accessKeyId</name>
<value>xxx</value>
</property>
<property>
<name>fs.oss.accessKeySecret</name>
<value>xxx</value>
</property>
</configuration>Use --disableChecksum
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
`--disableChecksum`: Disables file checksum verification.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --disableChecksumUse --deleteOnSuccess
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Supported |
`--deleteOnSuccess`: Moves data instead of copying it. This option is similar to an `mv` operation. It first copies the files and then deletes them from the source.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --deleteOnSuccessUse --enableTransaction
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Support |
`--enableTransaction`: By default, JindoDistCp ensures task-level integrity. You can use this parameter to ensure job-level integrity and enable transactional support between jobs.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --enableTransactionUse --ignore
Version | OSS | OSS-HDFS |
4.3.0 or later | Supported | Support |
`--ignore`: Ignores exceptions that occur during data migration. Errors do not interrupt the task. Instead, the errors are reported as JindoCounter values. If CMS is enabled, a notification is also sent.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --ignoreUse --diff
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
`--diff`: Compares the files in the source and destination. If a source file is not synchronized to the destination, a file that contains the differences is generated in the current directory. If your JindoDistCp task involves compression or decompression, `--diff` cannot show the correct file differences because the file size is changed during the process.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --diffIf differences are found, a file that contains the differences is generated in the current directory, and the following message is displayed:
JindoCounter
DIFF_FILES=1If your `--dest` is an HDFS path, the `/path`, `hdfs://hostname:ip/path`, and `hdfs://headerIp:ip/path` formats are supported. The `hdfs:///path`, `hdfs:/path`, or other custom formats are not supported.
To view differences in file metadata, run the --diff --preserveMeta command:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --diff --preserveMetaUse --update
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Support |
`--update`: Enables incremental synchronization. This option skips identical files and directories and synchronizes only new or changed files and directories from the source to the destination.
If a JindoDistCp task fails, you can use this parameter to resume it from the breakpoint and copy only the remaining files. You can also use this parameter to copy new files that are added to the source after the previous JindoDistCp task is completed.
The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --updateWrite data to OSS in Cold Archive, Archive, or IA storage classes
Version | OSS | OSS-HDFS |
4.3.0 or later | Support | Not supported |
`--policy`: Specifies the storage class for data that is written to OSS. You can set this parameter to Cold Archive, Archive, or IA. If you do not specify this parameter, data is written to the Standard storage class by default.
Write data to the Cold Archive storage class (coldArchive) in OSS
This feature is available only in specific regions. For more information, see OSS storage classes. The following command provides an example:
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-bucket/hourly_table --policy coldArchive --parallelism 20Write data to the Archive storage class (archive) in OSS
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-bucket/hourly_table --policy archive --parallelism 20Write data to the IA storage class (ia) in OSS
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-bucket/hourly_table --policy ia --parallelism 20
Use --preserveMeta
Version | OSS | OSS-HDFS |
4.4.0 or later | Not supported | Supported |
`--preserveMeta`: Specifies that metadata is migrated along with the data. The metadata includes Owner, Group, Permission, Atime, Mtime, Replication, BlockSize, XAttrs, and ACL.
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --preserveMetaUse --jobBatch
Version | OSS | OSS-HDFS |
4.5.1 or later | Support | Support |
`--jobBatch`: When your DistCp task writes data to OSS, you can use `--jobBatch` to specify the number of files to be processed by each DistCp job. The default value is 10,000.
jindo-distcp-tool-${version}.jar --src /data/hourly_table --dest oss://example-oss-bucket/hourly_table --jobBatch 50000Use --enableCMS
Version | OSS | OSS-HDFS |
4.5.1 or later | Support | Support |
`--enableCMS`: Enables the CMS alert feature.
JindoDistCp counters
JindoDistCp counters summarize the results of a JindoDistCp task. The following table describes the counters.
Parameter | Description |
COPY_FAILED | The number of files that failed to be copied. |
CHECKSUM_DIFF | The number of files that failed checksum verification. This is included in COPY_FAILED. |
FILES_EXPECTED | The number of files expected to be copied. |
BYTES_EXPECTED | The number of bytes expected to be copied. |
FILES_COPIED | The number of files that were successfully copied. |
BYTES_COPIED | The number of bytes that were successfully copied. |
FILES_SKIPPED | The number of files skipped during an incremental update. |
BYTES_SKIPPED | The number of bytes skipped during an incremental update. |
DIFF_FILES | The number of files that are different between the source and destination paths. |
SAME_FILES | The number of files that are identical between the source and destination paths. |
DST_MISS | The number of files that do not exist in the destination path. This is included in DIFF_FILES. |
LENGTH_DIFF | The number of files whose sizes are different between the source and destination. This is included in DIFF_FILES. |
CHECKSUM_DIFF | The number of files that failed checksum verification. This is included in DIFF_FILES. |
DIFF_FAILED | The number of files for which the comparison operation was abnormal. |