This tutorial provides a hands-on look at how you can move data out of Pulsar without writing a single line of code. It is helpful to review the concepts for Pulsar I/O in tandem with running the steps in this guide to gain a deeper understanding. At the end of this tutorial, you will be able to:

  • Connect your Pulsar cluster with your Cassandra cluster

Tip

  1. These instructions assume you are running Pulsar in standalone mode. However all the commands used in this tutorial should be able to be used in a multi-nodes Pulsar cluster without any changes.

  2. All the instructions are assumed to run at the root directory of a Pulsar binary distribution.

安装 Pulsar

开始运行Pulsar之前,请先用下列几种方式下载二进制包:

下载包后, 解压它和 cd 到生成的目录中:

  1. $ tar xvfz apache-pulsar-2.6.1-bin.tar.gz
  2. $ cd apache-pulsar-2.6.1

Installing Builtin Connectors

Since release 2.3.0, Pulsar releases all the builtin connectors as individual archives. If you would like to enable those builtin connectors, you can download the connectors “NAR” archives and from the Pulsar downloads page.

After downloading the desired builtin connectors, these archives should be places under the connectors directory where you have unpacked the Pulsar distribution.

  1. # Unpack regular Pulsar tarball and copy connectors NAR archives
  2. $ tar xvfz /path/to/apache-pulsar-2.6.1-bin.tar.gz
  3. $ cd apache-pulsar-2.6.1
  4. $ mkdir connectors
  5. $ cp -r /path/to/downloaded/connectors/*.nar ./connectors
  6. $ ls connectors
  7. pulsar-io-aerospike-2.6.1.nar
  8. pulsar-io-cassandra-2.6.1.nar
  9. pulsar-io-kafka-2.6.1.nar
  10. pulsar-io-kinesis-2.6.1.nar
  11. pulsar-io-rabbitmq-2.6.1.nar
  12. pulsar-io-twitter-2.6.1.nar
  13. ...

Tip

You can also use the Docker image apachepulsar/pulsar-all:2.6.1 which already comes with all the available builtin connectors.

Start Pulsar Service

  1. bin/pulsar standalone

All the components of a Pulsar service will start in order. You can curl those pulsar service endpoints to make sure Pulsar service is up running correctly.

  1. Check pulsar binary protocol port.
  1. telnet localhost 6650
  1. Check pulsar function cluster
  1. curl -s http://localhost:8080/admin/v2/worker/cluster

Example output:

  1. [{"workerId":"c-standalone-fw-localhost-6750","workerHostname":"localhost","port":6750}]
  1. Make sure public tenant and default namespace exist
  1. curl -s http://localhost:8080/admin/v2/namespaces/public

Example outoupt:

  1. ["public/default","public/functions"]
  1. All builtin connectors should be listed as available.
  1. curl -s http://localhost:8080/admin/v2/functions/connectors

Example output:

  1. [{"name":"aerospike","description":"Aerospike database sink","sinkClass":"org.apache.pulsar.io.aerospike.AerospikeStringSink"},{"name":"cassandra","description":"Writes data into Cassandra","sinkClass":"org.apache.pulsar.io.cassandra.CassandraStringSink"},{"name":"kafka","description":"Kafka source and sink connector","sourceClass":"org.apache.pulsar.io.kafka.KafkaStringSource","sinkClass":"org.apache.pulsar.io.kafka.KafkaBytesSink"},{"name":"kinesis","description":"Kinesis sink connector","sinkClass":"org.apache.pulsar.io.kinesis.KinesisSink"},{"name":"rabbitmq","description":"RabbitMQ source connector","sourceClass":"org.apache.pulsar.io.rabbitmq.RabbitMQSource"},{"name":"twitter","description":"Ingest data from Twitter firehose","sourceClass":"org.apache.pulsar.io.twitter.TwitterFireHose"}]

If an error occurred while starting Pulsar service, you may be able to seen exception at the terminal you are running pulsar/standalone, or you can navigate the logs directory under the Pulsar directory to view the logs.

Connect Pulsar to Apache Cassandra

Make sure you have docker available at your laptop. If you don’t have docker installed, you can follow the instructions.

We are using cassandra docker image to start a single-node cassandra cluster in Docker.

Setup the Cassandra Cluster

Start a Cassandra Cluster

  1. docker run -d --rm --name=cassandra -p 9042:9042 cassandra

Before moving to next steps, make sure the cassandra cluster is up running.

  1. Make sure the docker process is running.
  1. docker ps
  1. Check the cassandra logs to make sure cassandra process is running as expected.
  1. docker logs cassandra
  1. Check the cluster status
  1. docker exec cassandra nodetool status

Example output:

  1. Datacenter: datacenter1
  2. =======================
  3. Status=Up/Down
  4. |/ State=Normal/Leaving/Joining/Moving
  5. -- Address Load Tokens Owns (effective) Host ID Rack
  6. UN 172.17.0.2 103.67 KiB 256 100.0% af0e4b2f-84e0-4f0b-bb14-bd5f9070ff26 rack1

Create keyspace and table

We are using cqlsh to connect to the cassandra cluster to create keyspace and table.

  1. $ docker exec -ti cassandra cqlsh localhost
  2. Connected to Test Cluster at localhost:9042.
  3. [cqlsh 5.0.1 | Cassandra 3.11.2 | CQL spec 3.4.4 | Native protocol v4]
  4. Use HELP for help.
  5. cqlsh>

All the following commands are executed in cqlsh.

Create keyspace pulsar_test_keyspace
  1. cqlsh> CREATE KEYSPACE pulsar_test_keyspace WITH replication = {'class':'SimpleStrategy', 'replication_factor':1};

Create table pulsar_test_table

  1. cqlsh> USE pulsar_test_keyspace;
  2. cqlsh:pulsar_test_keyspace> CREATE TABLE pulsar_test_table (key text PRIMARY KEY, col text);

Configure a Cassandra Sink

Now that we have a Cassandra cluster running locally. In this section, we will configure a Cassandra sink connector. The Cassandra sink connector will read messages from a Pulsar topic and write the messages into a Cassandra table.

In order to run a Cassandra sink connector, you need to prepare a yaml config file including informations that Pulsar IO runtime needs to know. For example, how Pulsar IO can find the cassandra cluster, what is the keyspace and table that Pulsar IO will be using for writing Pulsar messages to.

Create a file examples/cassandra-sink.yml and edit it to fill in following content:

  1. configs:
  2. roots: "localhost:9042"
  3. keyspace: "pulsar_test_keyspace"
  4. columnFamily: "pulsar_test_table"
  5. keyname: "key"
  6. columnName: "col"

To learn more about Cassandra Connector, see Cassandra Connector.

Submit a Cassandra Sink

Pulsar provides the CLI for running and managing Pulsar I/O connectors.

We can run following command to sink a sink connector with type cassandra and config file examples/cassandra-sink.yml.

Note

The sink-type parameter of the currently built-in connectors is determined by the setting of the name parameter specified in the pulsar-io.yaml file.

  1. bin/pulsar-admin sink create \
  2. --tenant public \
  3. --namespace default \
  4. --name cassandra-test-sink \
  5. --sink-type cassandra \
  6. --sink-config-file examples/cassandra-sink.yml \
  7. --inputs test_cassandra

Once the command is executed, Pulsar will create a sink connector named cassandra-test-sink and the sink connector will be running as a Pulsar Function and write the messages produced in topic test_cassandra to Cassandra table pulsar_test_table.

Inspect the Cassandra Sink

You can use sink CLI and source CLI for inspecting and managing the IO connectors.

Retrieve Sink Info

  1. bin/pulsar-admin sink get \
  2. --tenant public \
  3. --namespace default \
  4. --name cassandra-test-sink

Example output:

  1. {
  2. "tenant": "public",
  3. "namespace": "default",
  4. "name": "cassandra-test-sink",
  5. "className": "org.apache.pulsar.io.cassandra.CassandraStringSink",
  6. "inputSpecs": {
  7. "test_cassandra": {
  8. "isRegexPattern": false
  9. }
  10. },
  11. "configs": {
  12. "roots": "localhost:9042",
  13. "keyspace": "pulsar_test_keyspace",
  14. "columnFamily": "pulsar_test_table",
  15. "keyname": "key",
  16. "columnName": "col"
  17. },
  18. "parallelism": 1,
  19. "processingGuarantees": "ATLEAST_ONCE",
  20. "retainOrdering": false,
  21. "autoAck": true,
  22. "archive": "builtin://cassandra"
  23. }

Check Sink Running Status

  1. bin/pulsar-admin sink status \
  2. --tenant public \
  3. --namespace default \
  4. --name cassandra-test-sink

Example output:

  1. {
  2. "numInstances" : 1,
  3. "numRunning" : 1,
  4. "instances" : [ {
  5. "instanceId" : 0,
  6. "status" : {
  7. "running" : true,
  8. "error" : "",
  9. "numRestarts" : 0,
  10. "numReadFromPulsar" : 0,
  11. "numSystemExceptions" : 0,
  12. "latestSystemExceptions" : [ ],
  13. "numSinkExceptions" : 0,
  14. "latestSinkExceptions" : [ ],
  15. "numWrittenToSink" : 0,
  16. "lastReceivedTime" : 0,
  17. "workerId" : "c-standalone-fw-localhost-8080"
  18. }
  19. } ]
  20. }

Verify the Cassandra Sink

Now lets produce some messages to the input topic of the Cassandra sink test_cassandra.

  1. for i in {0..9}; do bin/pulsar-client produce -m "key-$i" -n 1 test_cassandra; done

Inspect the sink running status again. You should be able to see 10 messages are processed by the Cassandra sink.

  1. bin/pulsar-admin sink status \
  2. --tenant public \
  3. --namespace default \
  4. --name cassandra-test-sink

Example output:

  1. {
  2. "numInstances" : 1,
  3. "numRunning" : 1,
  4. "instances" : [ {
  5. "instanceId" : 0,
  6. "status" : {
  7. "running" : true,
  8. "error" : "",
  9. "numRestarts" : 0,
  10. "numReadFromPulsar" : 10,
  11. "numSystemExceptions" : 0,
  12. "latestSystemExceptions" : [ ],
  13. "numSinkExceptions" : 0,
  14. "latestSinkExceptions" : [ ],
  15. "numWrittenToSink" : 10,
  16. "lastReceivedTime" : 1551685489136,
  17. "workerId" : "c-standalone-fw-localhost-8080"
  18. }
  19. } ]
  20. }

Finally, lets inspect the results in Cassandra using cqlsh

  1. docker exec -ti cassandra cqlsh localhost

Select the rows from the Cassandra table pulsar_test_table:

  1. cqlsh> use pulsar_test_keyspace;
  2. cqlsh:pulsar_test_keyspace> select * from pulsar_test_table;
  3. key | col
  4. --------+--------
  5. key-5 | key-5
  6. key-0 | key-0
  7. key-9 | key-9
  8. key-2 | key-2
  9. key-1 | key-1
  10. key-3 | key-3
  11. key-6 | key-6
  12. key-7 | key-7
  13. key-4 | key-4
  14. key-8 | key-8

Delete the Cassandra Sink

  1. bin/pulsar-admin sink delete \
  2. --tenant public \
  3. --namespace default \
  4. --name cassandra-test-sink