Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 3 additions & 6 deletions modules/ROOT/pages/astream-subscriptions-exclusive.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -25,15 +25,13 @@ To create a {pulsar-short} exclusive subscription, create a `pulsarConsumer` wit

. In `src/main/java/com/datastax/pulsar`, create a `SimplePulsarConsumer.java` file with the following contents:
+
.SimplePulsarConsumer.java
[source,java,subs="+attributes"]
----
include::ROOT:partial$simplepulsarconsumer.java[]
----
+
Alternatively, you can omit the `.subscriptionType` declaration because `Exclusive` is the default subscription type:
+
.SimplePulsarConsumer.java with implied exclusive subscription
[source,java]
----
pulsarConsumer = pulsarClient.newConsumer(Schema.JSON(DemoBean.class))
Expand All @@ -52,9 +50,10 @@ Alternatively, you can omit the `.subscriptionType` declaration because `Exclusi

include::ROOT:partial$subscription-start-consumer.adoc[]

. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages:
. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages.
+
The producer's terminal shows when each message is sent:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarProducer - Message 93573631 sent
Expand All @@ -64,7 +63,6 @@ include::ROOT:partial$subscription-start-consumer.adoc[]
+
In the `SimplePulsarConsumer` terminal, the consumer begins consuming the produced messages:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarConsumer - Message received: {"show_id":93573631,"cast":"LeBron James, Anthony Davis, Kyrie Irving, Damian Lillard, Klay Thompson...","country":"United States","date_added":"July 16, 2021","description":"NBA superstar LeBron James teams up with Bugs Bunny and the rest of the Looney Tunes for this long-awaited sequel.","director":"Malcolm D. Lee","duration":"120 min","listed_in":"Animation, Adventure, Comedy","rating":"PG","release_year":2021,"title":"Space Jam: A New Legacy","type":"Movie"}
Expand All @@ -76,7 +74,6 @@ In the `SimplePulsarConsumer` terminal, the consumer begins consuming the produc
+
The second consumer cannot subscribe to the topic because the subscription is exclusive:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.Configuration - Configuration has been loaded successfully
Expand Down
8 changes: 3 additions & 5 deletions modules/ROOT/pages/astream-subscriptions-failover.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@ To create a {pulsar-short} failover subscription, create a `pulsarConsumer` with

. In `src/main/java/com/datastax/pulsar`, create a `SimplePulsarConsumer.java` file with the following contents:
+
.SimplePulsarConsumer.java
[source,java,subs="+attributes"]
----
include::ROOT:partial$simplepulsarconsumer.java[]
Expand All @@ -37,9 +36,10 @@ include::ROOT:partial$simplepulsarconsumer.java[]

include::ROOT:partial$subscription-start-consumer.adoc[]

. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages:
. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages.
+
The producer's terminal shows when each message is sent:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarProducer - Message 50585599 sent
Expand All @@ -52,7 +52,6 @@ include::ROOT:partial$subscription-start-consumer.adoc[]
+
In the `SimplePulsarConsumer` terminal, the primary consumer begins consuming messages:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarConsumer - Message received: {"show_id":50585599,"cast":"LeBron James, Anthony Davis, Kyrie Irving, Damian Lillard, Klay Thompson...","country":"United States","date_added":"July 16, 2021","description":"NBA superstar LeBron James teams up with Bugs Bunny and the rest of the Looney Tunes for this long-awaited sequel.","director":"Malcolm D. Lee","duration":"120 min","listed_in":"Animation, Adventure, Comedy","rating":"PG","release_year":2021,"title":"Space Jam: A New Legacy","type":"Movie"}
Expand All @@ -66,7 +65,6 @@ The backup consumer subscribes to the topic but does not immediately begin consu
. In your first `SimplePulsarConsumer` terminal, stop the process (`Ctrl+C`), and then switch to your second `SimplePulsarConsumer` terminal.
Notice that the backup consumer begins consuming messages where the first consumer left off:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarConsumer - Message received: {"show_id":73260535,"cast":"LeBron James, Anthony Davis, Kyrie Irving, Damian Lillard, Klay Thompson...","country":"United States","date_added":"July 16, 2021","description":"NBA superstar LeBron James teams up with Bugs Bunny and the rest of the Looney Tunes for this long-awaited sequel.","director":"Malcolm D. Lee","duration":"120 min","listed_in":"Animation, Adventure, Comedy","rating":"PG","release_year":2021,"title":"Space Jam: A New Legacy","type":"Movie"}
Expand Down
15 changes: 3 additions & 12 deletions modules/ROOT/pages/astream-subscriptions-keyshared.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -42,15 +42,13 @@ Running multiple consumers with `autoSplitHashRange` balances the messaging load

. In `src/main/java/com/datastax/pulsar`, create a `SimplePulsarConsumer.java` file with the following contents:
+
.SimplePulsarConsumer.java
[source,java,subs="+attributes"]
----
include::ROOT:partial$simplepulsarconsumer.java[]
----

. In `pulsarConsumer`, add `.keySharedPolicy(KeySharedPolicy.autoSplitHashRange())`:
+
.SimplePulsarConsumer.java with autoSplitHashRange
[source,java,subs="+attributes"]
----
...
Expand Down Expand Up @@ -80,15 +78,13 @@ This policy requires additional dependencies and producer configuration changes,

. In `src/main/java/com/datastax/pulsar`, create a `SimplePulsarConsumer.java` file with the following contents:
+
.SimplePulsarConsumer.java
[source,java,subs="+attributes"]
----
include::ROOT:partial$simplepulsarconsumer.java[]
----

. Import the following additional classes that are required for the `stickyHashRange` policy:
+
.SimplePulsarConsumer.java
[source,java]
----
import org.apache.pulsar.client.api.Range;
Expand All @@ -100,7 +96,6 @@ import org.apache.pulsar.client.api.SubscriptionType;
+
The following example sets all possible hashes (`0-65535`) on this subscription to one consumer:
+
.SimplePulsarConsumer.java with stickyHashRange for one consumer
[source,java,subs="+attributes"]
----
...
Expand All @@ -125,7 +120,6 @@ pulsarConsumer = pulsarClient.newConsumer(Schema.JSON(DemoBean.class))
To split the hash range between multiple consumers, add a `Range.of()` argument for each consumer with the assigned hash range.
For example:
+
.SimplePulsarConsumer.java with stickyHashRange for two consumers
[source,java]
----
// Policy assigns half of the hash range to one consumer and half to another
Expand All @@ -136,7 +130,6 @@ For example:
+
.. In `SimplePulsarProducer.java`, import the following classes:
+
.SimplePulsarProducer.java
[source,java]
----
import org.apache.pulsar.client.api.BatcherBuilder;
Expand All @@ -145,7 +138,6 @@ import org.apache.pulsar.client.api.HashingScheme;

.. Configure the `pulsarProducer` to use the `JavaStringHash` hashing scheme:
+
.SimplePulsarProducer.java
[source,java]
----
pulsarProducer = pulsarClient
Expand All @@ -165,9 +157,10 @@ pulsarProducer = pulsarClient

include::ROOT:partial$subscription-start-consumer.adoc[]

. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages:
. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages.
+
The producer's terminal shows when each message is sent:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarProducer - Message 55794190 sent
Expand All @@ -178,7 +171,6 @@ include::ROOT:partial$subscription-start-consumer.adoc[]
+
In the `SimplePulsarConsumer` terminal, the consumer begins receiving messages:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarConsumer - Message received: {"show_id":55794190,"cast":"LeBron James, Anthony Davis, Kyrie Irving, Damian Lillard, Klay Thompson...","country":"United States","date_added":"July 16, 2021","description":"NBA superstar LeBron James teams up with Bugs Bunny and the rest of the Looney Tunes for this long-awaited sequel.","director":"Malcolm D. Lee","duration":"120 min","listed_in":"Animation, Adventure, Comedy","rating":"PG","release_year":2021,"title":"Space Jam: A New Legacy","type":"Movie"}
Expand All @@ -193,7 +185,6 @@ The auto-hashing policy balances hash ranges across available consumers.
If you used sticky hashing with one `Range.of()` argument, then the new consumer cannot subscribe to the topic because the `SimplePulsarConsumer` configuration reserved the entire hash range for the first consumer.
For example:
+
.Result when using sticky hashing limited to one consumer
[source,console]
----
[main] INFO com.datastax.pulsar.Configuration - Configuration has been loaded successfully
Expand Down
8 changes: 3 additions & 5 deletions modules/ROOT/pages/astream-subscriptions-shared.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ To create a {pulsar-short} shared subscription, create a `pulsarConsumer` with `

. In `src/main/java/com/datastax/pulsar`, create a `SimplePulsarConsumer.java` file with the following contents:
+
.SimplePulsarConsumer.java
[source,java,subs="+attributes"]
----
include::ROOT:partial$simplepulsarconsumer.java[]
Expand All @@ -36,9 +35,10 @@ include::ROOT:partial$simplepulsarconsumer.java[]

include::ROOT:partial$subscription-start-consumer.adoc[]

. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages:
. In a new terminal window, run `SimplePulsarProducer.java` to begin producing messages.
+
The producer's terminal shows when each message is sent:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarProducer - Message 59819331 sent
Expand All @@ -50,7 +50,6 @@ include::ROOT:partial$subscription-start-consumer.adoc[]
+
In the `SimplePulsarConsumer` terminal, the consumer begins receiving messages:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarConsumer - Message received: {"show_id":59819331,"cast":"LeBron James, Anthony Davis, Kyrie Irving, Damian Lillard, Klay Thompson...","country":"United States","date_added":"July 16, 2021","description":"NBA superstar LeBron James teams up with Bugs Bunny and the rest of the Looney Tunes for this long-awaited sequel.","director":"Malcolm D. Lee","duration":"120 min","listed_in":"Animation, Adventure, Comedy","rating":"PG","release_year":2021,"title":"Space Jam: A New Legacy","type":"Movie"}
Expand All @@ -62,7 +61,6 @@ In the `SimplePulsarConsumer` terminal, the consumer begins receiving messages:
+
The new consumer subscribes to the topic and consumes messages:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.SimplePulsarConsumer - Message received: {"show_id":70129519,"cast":"LeBron James, Anthony Davis, Kyrie Irving, Damian Lillard, Klay Thompson...","country":"United States","date_added":"July 16, 2021","description":"NBA superstar LeBron James teams up with Bugs Bunny and the rest of the Looney Tunes for this long-awaited sequel.","director":"Malcolm D. Lee","duration":"120 min","listed_in":"Animation, Adventure, Comedy","rating":"PG","release_year":2021,"title":"Space Jam: A New Legacy","type":"Movie"}
Expand Down
6 changes: 4 additions & 2 deletions modules/ROOT/partials/sinks/edit.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@ Additionally, some properties can be modified with specific arguments, such as `

To get the current configuration, see xref:apis:api-operations.adoc#get-sink-connector-configuration-data[Get sink connector configuration data].

.pulsar-admin CLI
pulsar-admin CLI::
+
[source,shell,subs="+attributes"]
----
./bin/pulsar-admin sinks update \
Expand All @@ -17,7 +18,8 @@ To get the current configuration, see xref:apis:api-operations.adoc#get-sink-con
--parallelism 2
----

.{pulsar-short} Admin API
{pulsar-short} Admin API::
+
[source,shell]
----
curl -sS --fail -L -X PUT "$WEB_SERVICE_URL/admin/v3/sinks/$TENANT/$NAMESPACE/$SINK_NAME" \
Expand Down
9 changes: 6 additions & 3 deletions modules/ROOT/partials/sinks/get-started.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ For example: `{connectorType}-sink-prod-us-east-1`.
. Create the connector using JSON-formatted connector configuration data.
You can pass the configuration directly or with a configuration file.
+
.pulsar-admin CLI
pulsar-admin CLI::
+
[source,shell,subs="+attributes"]
----
./bin/pulsar-admin sinks create \
Expand All @@ -28,15 +29,17 @@ You can pass the configuration directly or with a configuration file.
--sink-config-file configs.json
----
+
.{pulsar-short} Admin API
{pulsar-short} Admin API::
+
[source,shell]
----
curl -sS --fail -L -X POST "$WEB_SERVICE_URL/admin/v3/sinks/$TENANT/$NAMESPACE/$SINK_NAME" \
--header "Authorization: Bearer $PULSAR_TOKEN" \
--form "sinkConfig=@configs.json;type=application/json"
----
+
.Example configuration data structure
Example configuration data structure::
+
[source,json]
----
include::common:streaming:example$connectors/sinks/{connectorType}/sample-data.json[]
Expand Down
6 changes: 4 additions & 2 deletions modules/ROOT/partials/sources/edit.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@ Additionally, some properties can be modified with specific arguments, such as `

To get the current configuration, see xref:apis:api-operations.adoc#get-source-connector-configuration-data[Get source connector configuration data].

.pulsar-admin CLI
pulsar-admin CLI::
+
[source,shell,subs="+attributes"]
----
./bin/pulsar-admin sources update \
Expand All @@ -17,7 +18,8 @@ To get the current configuration, see xref:apis:api-operations.adoc#get-source-c
--parallelism 2
----

.{pulsar-short} Admin API
{pulsar-short} Admin API::
+
[source,shell]
----
curl -sS --fail -L -X PUT "$WEB_SERVICE_URL/admin/v3/sources/$TENANT/$NAMESPACE/$SOURCE_NAME" \
Expand Down
9 changes: 6 additions & 3 deletions modules/ROOT/partials/sources/get-started.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ For example: `{connectorType}-source-prod-us-east-1`.
. Create the connector using JSON-formatted connector configuration data.
You can pass the configuration directly or with a configuration file.
+
.pulsar-admin CLI
pulsar-admin CLI::
+
[source,shell,subs="+attributes"]
----
./bin/pulsar-admin sources create \
Expand All @@ -28,15 +29,17 @@ You can pass the configuration directly or with a configuration file.
--source-config-file configs.json
----
+
.{pulsar-short} Admin API
{pulsar-short} Admin API::
+
[source,shell]
----
curl -sS --fail -L -X POST "$WEB_SERVICE_URL/admin/v3/sources/$TENANT/$NAMESPACE/$SOURCE_NAME" \
--header "Authorization: Bearer $PULSAR_TOKEN" \
--form "sourceConfig=@mynetty-source-config.json;type=application/json"
----
+
.Example configuration data structure
Example configuration data structure::
+
[source,json]
----
include::common:streaming:example$connectors/sources/{connectorType}/sample-data.json[]
Expand Down
17 changes: 7 additions & 10 deletions modules/ROOT/partials/subscription-setup-project.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@

. Edit the `pom.xml` file to include the following dependencies:
+
.pom.xml
[source,xml]
----
<project xmlns="http://maven.apache.org/POM/4.0.0"
Expand Down Expand Up @@ -73,11 +72,11 @@
</project>
----

. In `src/main/resources`, create the following `application.properties` file with the connection details for your {product} cluster.
. In `src/main/resources`, create an `application.properties` file with the connection details for your {product} cluster.
+
Create the `resources` subdirectory if it doesn't already exist.
+
./src/main/resources/application.properties
[source,properties,subs="+quotes"]
[source,plaintext,subs="+quotes"]
----
# ---------------------------------------
# Configuration of your Astra Streaming tenant
Expand All @@ -91,10 +90,10 @@ authentication_token=**ASTRA_APPLICATION_TOKEN**
topic_name=my-topic
----

. In `src/main/java/com/datastax/pulsar`, create the following `Configuration.java` class to load the connection details from `application.properties`.
. In `src/main/java/com/datastax/pulsar`, create a `Configuration.java` class to load the connection details from `application.properties`.
+
Create the `/datastax/pulsar` subdirectories if they don't already exist.
+
.Configuration.java
[source,java]
----
package com.datastax.pulsar;
Expand Down Expand Up @@ -201,9 +200,8 @@ public class Configuration {
}
----

. In `src/main/java/com/datastax/pulsar`, create the following `DemoBean.java` class to represent the example messages that will be produced and consumed:
. In `src/main/java/com/datastax/pulsar`, create a `DemoBean.java` class to represent the example messages that will be produced and consumed:
+
.DemoBean.java
[source,java]
----
package com.datastax.pulsar;
Expand Down Expand Up @@ -349,9 +347,8 @@ public class DemoBean {
}
----

. In `src/main/java/com/datastax/pulsar`, create a `SimplePulsarProducer.java` file with the following contents:
. In `src/main/java/com/datastax/pulsar`, create a file named `SimplePulsarProducer.java` with the following contents:
+
.SimplePulsarProducer.java
[source,java]
----
package com.datastax.pulsar;
Expand Down
3 changes: 1 addition & 2 deletions modules/ROOT/partials/subscription-start-consumer.adoc
Original file line number Diff line number Diff line change
@@ -1,8 +1,7 @@
. Run `SimplePulsarConsumer.java` to begin consuming messages as the primary consumer.
+
The confirmation message and a cursor appear to indicate the consumer is ready:
A confirmation message and a cursor appear when the consumer is ready:
+
.Result
[source,console]
----
[main] INFO com.datastax.pulsar.Configuration - Configuration has been loaded successfully
Expand Down
Loading