Skip to main content

pydistcp: python WebHDFS inter/intra-cluster data copy tool.

Project description

pypi

pydistcp

A python WebHDFS/HTTPFS based tool for inter/intra-cluster data copying. This tool is very suitable for multiple mid or small size files cross-clusters copy. Compared to the normal distcp which adds a lot of overhead time for submitting the map-reduce job then waiting for YARN to schedule it..., pydistcp uses webhdfs to stream the data from source cluster datanodes directly to destination cluster datanodes using multiple parallel threads.

When transferring few huge files, the normal distcp may be faster, but when transferring lot of small, midsize or relatively big file, pydistcp provides a very good performance.

  $ pydistcp -f -s staging -d prod /data/outgoing /data/incoming --threads=10 --part-size=131072
  27.1%   [ pending: 32 | transferring: 6 | complete: 4 ]
Job Status:
{
  "Size Failed": 0,
  "Size Copied": 257721641,
  "Source Path": "/data/t100",
  "Size Expected": 257721641,
  "Files Expected": 42,
  "Files Failed": 0,
  "Destination Path": "/data/t200",
  "Start Time": "2017-02-22 17:39:29",
  "Files Skipped": 0,
  "Size Deleted": 0,
  "End Time": "2017-02-22 17:39:50",
  "Files Copied": 42,
  "Files Deleted": 0,
  "Duration": 20.756325006484985,
  "Outcome": "Successful",
  "Size Skipped": 0
}

Pydistcp uses pywhdfs for establishing connections with WEBHDFS/HTTPFS source and destination clusters.

Features

  • Pydistcp is based on pywhdfs project to establish WebHDFS and HTTPFS connections with source and destination clusters, so all clusters configurations supported in pywhdfs are also supported in pydistcp:
    • Support both secure (Kerberos,Token) and insecure clusters
    • Supports HA cluster and handle namenode failover
    • Supports HDFS federation with multiple nameservices and mount points.
  • Supports data copy between secure and insecure clusters
  • Supports data copy between clusters using different kerberos realms using token authentication
  • Supports data copy between encrypted and unencrypted clusters
  • Json format clusters configuration.
  • Perform concurrent multithreaded data copy.

Getting started

  $ easy_install pydistcp

Configuration

Pydistcp share the same json configuration file used by pywhdfs . Please refer to the project readme file for details about the json configuration schema.

USAGE

There are multiple arguments you can use to alter the way the copy works, or to enhance the performance of the job depending on the size of the server you use. Use the help argument to display the full list of supported parameters:

  $ pydistcp --help
  pydistcp: A python Web HDFS based tool for inter/intra-cluster data copying.

  Usage:
    pydistcp [-fp] [--no-checksum] [--silent] (-s CLUSTER -d CLUSTER) [-v...] [--part-size=PART_SIZE] [--threads=THREADS] SRC_PATH DEST_PATH
    pydistcp (--version | -h)

  Options:
    --version                     Show version and exit.
    -h --help                     Show help and exit.
    -s CLUSTER --src=CLUSTER      Alias of source namenode to connect to (valid only with dist).
    -d CLUSTER --dest=CLUSTER     Alias of destination namenode to connect to (valid only with dist).
    -v --verbose                  Enable log output. Can be specified multiple times to increase verbosity each time.
    --no-checksum                 Disable checksum check prior to file transfer. This will force overwrite.
    --silent                      Don't display progress status.
    -f --force                    Allow overwriting any existing files.
    -p --preserve                 Preserve file attributes.
    --threads=THREADS             Number of threads to use for parallelization.
                                  zero limits the concurrency to the maximum concurrent threads
                                  supported by the cluster. [default: 0]
    --part-size=PART_SIZE         Interval in bytes by which the files will be copied
                                  needs to be a Powers of 2. [default: 65536]

  Examples:
    pydistcp -s prod -d preprod -v /tmp/src /tmp/dest

All cluster connection parameters will be fetched from the json configuration file.

benchmarks

Below some benchmarks showing the impact of data size on the copy performance using pydistcp :

File Count Data Size Time
2379 11.4 G 4m39.069s
242 25.9 G 5m39.348s
869 116.9 G 25m53.231s
42 545.8 M 0m19.946s
1788 5.2 G 2m25.649s
4428 35.7 G 10m20.129s
2357 5.6 G 3m2.598s
180 2.3 G 0m33.133s
334 7.6 G 1m26.260s

Note that all test cases are executed with 10 concurrent threads on a machine having 6 cores and supporting up to 12 threads and no files are skipped during the copy. Both the source and destination clusters are secured with kerberos and use ssl to encrypt transferred data.

Pydistcp performance may be impact by lot of parameters like:

  • the size of the machine performing the copy.
  • The type of the source and destination clusters (secure clusters with kerberos does not support lot of concurrent threads, it is better from a performance perspective to use token authentication)
  • SSL and the length of encryption key used
  • The type of data to be transferred : Pydistcp deliver the best performance for multiple files having approximately uniform sizes.

Contributing

Feedback and Pull requests are very welcome!

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

pydistcp-1.0.7.tar.gz (12.7 kB view details)

Uploaded Source

File details

Details for the file pydistcp-1.0.7.tar.gz.

File metadata

  • Download URL: pydistcp-1.0.7.tar.gz
  • Upload date:
  • Size: 12.7 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/1.15.0 pkginfo/1.5.0.1 requests/2.24.0 setuptools/44.1.1 requests-toolbelt/0.9.1 tqdm/4.48.2 CPython/2.7.10

File hashes

Hashes for pydistcp-1.0.7.tar.gz
Algorithm Hash digest
SHA256 5165ce517aa2b5e3c4d1b769ae0ae0adf3ca79c8587f806dbe7f153df3f1363d
MD5 4d2ed8e86942bde002fe77af8b35f6a5
BLAKE2b-256 9736a91ab002e4fbd0817fdd1e131424f533cadc9aacbcc56df6c1d4b55b775f

See more details on using hashes here.

Supported by

AWS AWS Cloud computing and Security Sponsor Datadog Datadog Monitoring Fastly Fastly CDN Google Google Download Analytics Microsoft Microsoft PSF Sponsor Pingdom Pingdom Monitoring Sentry Sentry Error logging StatusPage StatusPage Status page