Providing a way of migrating data in Kafka topics into tables in Hbase, preserving versions based on Kafka message timestamps.
Two columns are written to for each message received; one to store the body
of the message and one to store a count and last received date of the
topic. These are configured using the K2HB_KAFKA_TOPIC_* and
K2HB_KAFKA_DATA_* environment variables.
By default the data table is k2hb:ingest with a column family of topic.
The qualifier is the topic name, the body of the cell is the raw message
received from Kafka and the version is the timestamp of the message in
milliseconds.
Along with the data of the message a counter is kept for each topic to
indicate how many messages have been processed and when. This is useful for
creating a list of topics to process or limiting that list to only topics
that have received new data since a given time. The default table is
k2hb:ingest-topic and the default column is c:msg.
For example, after receiving a single message on test-topic the data
is as follows:
hbase(main):001:0> scan 'k2hb:ingest'
ROW COLUMN+CELL
63213667-c5a5-4411-a93b-e2da709c553e column=topic:test-topic, timestamp=1563547895682, value=<message body>
1 row(s) in 0.1090 seconds
hbase(main):002:0> scan 'k2hb:ingest-topic'
ROW COLUMN+CELL
test-topic column=c:msg, timestamp=1563547895689, value=\x00\x00\x00\x00\x00\x00\x00\x01
1 row(s) in 0.0100 seconds
Kafka2Hbase will attempt to create the required namespaces, tables and column families on startup. If they already exist, nothing will happen. By default the data table column family has a maximum of MAXINT versions (approximately 2.1 billion) and a minimum of 1 version. There is no TTL. The topic counter column family has no versioning or TTL.
A Makefile wraps some of the gradle and docker-compose commands to give a more unified basic set of operations. These can be checked by running:
$ make help
Ensure a JVM is installed and run the gradle wrapper.
make build
If a standard zip file is required, just use the assembleDist command.
make dist
This produces a zip and a tarball of the latest version.
A full local stack can be run using the provided Dockerfile and Docker Compose configuration. The Dockerfile uses a multi-stage build so no pre-compilation is required.
make up
The environment can be stopped without losing any data:
make down
Or completely removed including all data volumes:
make destroy
Integration tests can be executed inside a Docker container to make use of
the Kafka and Hbase instances running in the local stack. The integration
tests are written in Kotlin and use the standard kotlintest testing framework.
make integration-all -> to run from a clean build
make integration -> to run just the tests again with everything running
The unit tests use JUnit to run and are written using specification language. They can be executed with the following command.
make test
Both Kafka2HBase and the integration tests can be run in an IDE to facilitate quicker feedback then a containerized approach. This is useful during active development.
To do this first bring up the hbase, kafka and zookeeper containers:
make services
On the run configuration for Kafka2Hbase set the following environment variables (nb not system properties)
K2HB_HBASE_ZOOKEEPER_QUORUM=localhost;K2HB_KAFKA_POLL_TIMEOUT=PT2S
And on the run configuration for the integration tests set these:
K2HB_KAFKA_BOOTSTRAP_SERVERS=localhost:9092;K2HB_HBASE_ZOOKEEPER_QUORUM=localhost
Then insert into your local hosts file the names, IP addresses of the kafka and hbase containers:
./hosts.sh
The services are listed in the docker-compose.yaml file and logs can be
retrieved for all services, or for a subset.
docker-compose logs hbase
The logs can be followed so new lines are automatically shown.
docker-compose logs -f hbase
To access the HBase shell it's necessary to use a Docker container. This can be run as a separate container.
make hbase-shell
There are a number of environment variables that can be used to configure the system. Some of them are for configuring Kafka2Hbase itself, and some are for configuring the built-in ACM PCA client to perform mutual auth.
By default Kafka2Hbase will connect to Zookeeper at zookeeper:2181 use the parent uri hbase
and create tables in the k2hb namespace. The data will be stored in cf:data
with at least 1 version and at most 10 versions and a TTL of 10 days.
/hbase but should be set to /hbase-unsecure for AWS HBaseBy default Kafka2Hbase will connect to Kafka at kafka:9092 in the k2hb
consumer group. It will poll the test-topic topic with a poll timeout of
10 days, and refresh the topics list every 10 seconds (10000 ms).
db.*. Defaults to test-topic.*10000 ms (10 seconds).
Typically, should be an order of magnitude less than K2HB_KAFKA_POLL_TIMEOUT, else new topics will not be discovered within each polling interval.PT10S).
Defaults to 1 Hour.
Should be greater than K2HB_KAFKA_META_REFRESH_MS, else new topics will not be discovered within each polling interval.K2HB_KAFKA_INSECURE=trueCERTGEN or retrieve
them from ACM with value RETRIEVEBy default the SSL is enabled but has no defaults. These must either be
configured in full or disabled entirely via K2HB_KAFKA_INSECURE=FALSE
and K2HB_KAFKA_CERT_MODE=CERTGEN.
For an authoritative full list of arguments see the tool help; Arguments not listed here are
defaulted in the entrypoint.sh script.
RSA or DSA)1024, 2048 or 4096)sha256, sha384, sha512)SHA256WITHECDSA, SHA384WITHECDSA, SHA512WITHECDSA, SHA256WITHRSA, SHA384WITHRSA, SHA512WITHRSA)1y2m6d)CRITICAL, ERROR, WARNING, INFO, DEBUG)By default the SSL is enabled but has no defaults. These must either be
configured in full or disabled entirely via K2HB_KAFKA_INSECURE=FALSE
and K2HB_KAFKA_CERT_MODE=RETRIEVE.
For an authoritative full list of arguments see the tool help; Arguments not listed here are
defaulted in the entrypoint.sh script.
true, false, yes, no, 1 or 0
If missing defaults to falseCRITICAL, ERROR, WARNING, INFO, DEBUG)Content type
Image
Digest
Size
601.4 MB
Last updated
almost 6 years ago
docker pull dwpdigital/kafka-to-hbase