Configuring Connection Details

You configure the Solace Connector for Apache Spark entirely through .option() calls on one of the following Spark objects: a DataStreamReader (for streaming reads), a DataStreamWriter (for streaming writes), or a DataFrameWriter (for batch writes, for example inside forEachBatch). There is no separate configuration file. Before you configure the connection details, ensure you meet the prerequisites. For more information, see Before You Begin.

The Connector for Apache Spark supports both directions:

  • Solace to Apache Spark (source, streaming reads only)

  • Apache Spark to Solace (sink)

Connection and Authentication

The following options are common to both the source and the sink directions.

Config Option Type Valid Values Default Value Description

host

String

tcp://<host>:<port>

None

Required. The event broker host and port. The default Solace Message Format (SMF) port is 55555 for plaintext, for example tcp://mybroker:55555, or 55443 for TLS, for example tcps://mybroker:55443.

vpn

String

None

Required. The Message VPN to connect to.

username

String

None

Required unless you use client certificate or OAuth authentication. The client username.

password

String

None

Required unless you use client certificate or OAuth authentication. The client password. Treat this value as a secret.

connectRetries

int

Integer or -1

0

Optional. The number of times to retry the initial connection. Use -1 to retry forever.

reconnectRetries

int

3

Optional. The number of times to attempt reconnection after the connection to the event broker is lost.

connectRetriesPerHost

int

0

Optional. The number of connection retries per host, when host specifies a list of hosts. Applies together with connectRetries, not instead of it.

reconnectRetryWaitInMillis

int

0-60000

3000

Optional. The wait time, in milliseconds, between connection and reconnection attempts.

In addition to the preceding options, you can pass any raw JCSMP session property using the solace.apiProperties.<Property> prefix, for example solace.apiProperties.SSL_TRUST_STORE, solace.apiProperties.reapply_subscriptions=false, or solace.apiProperties.client_channel_properties.keepAliveIntervalInMillis=3000. The connector strips the prefix and applies the value directly to the underlying JCSMPProperties object.

Authentication Schemes

Set solace.apiProperties.AUTHENTICATION_SCHEME to select the authentication scheme.

Scheme Description

Basic (default)

Uses the username and password options directly. Active when AUTHENTICATION_SCHEME is either unset or set to AUTHENTICATION_SCHEME_BASIC.

Client Certificate

Set solace.apiProperties.AUTHENTICATION_SCHEME=AUTHENTICATION_SCHEME_CLIENT_CERTIFICATE. You must also configure the following standard JCSMP TLS passthrough options, each prefixed with solace.apiProperties.:

  • SSL_TRUST_STORE

  • SSL_TRUST_STORE_FORMAT

  • SSL_TRUST_STORE_PASSWORD

  • SSL_KEY_STORE

  • SSL_KEY_STORE_FORMAT

  • SSL_KEY_STORE_PASSWORD

OAuth 2.0 (Client Credentials)

Set solace.apiProperties.AUTHENTICATION_SCHEME=AUTHENTICATION_SCHEME_OAUTH2

The connector can either fetch a token from an authorization server (server-fetch mode) or read a token that another process rotates onto disk (file mode). For details about how to configure these modes, see:

OAuth Server-Fetch Mode

The following configuration options are available for server-fetch mode:

Config Option Type Default Value Description

solace.oauth.client.auth-server-url

String

None

Required for server-fetch mode. The OAuth token endpoint URL.

solace.oauth.client.client-id

String

None

Required for server-fetch mode. The OAuth client ID.

solace.oauth.client.credentials.client-secret

String

None

Required for server-fetch mode. The OAuth client secret. Treat this value as a secret.

solace.oauth.client.auth-server.client-certificate.file

String

None

Optional. Path to an X.509 client certificate used only when calling the OAuth token endpoint over mutual TLS (mTLS). This is independent of the event broker session TLS settings.

solace.oauth.client.auth-server.truststore.file

String

None

Optional. Path to the truststore used for the token endpoint call. If an existing JKS truststore is available, point at it directly. If auth-server.client-certificate.file is also set, provide a path that includes a file name. The connector loads the certificate into a keystore and writes the JKS to that path.

solace.oauth.client.auth-server.truststore.password

String

None

Required if you configure a truststore or client certificate for the token endpoint call.

solace.oauth.client.auth-server.truststore.type

String

JKS

Optional. The truststore type for the token endpoint call.

solace.oauth.client.auth-server.ssl.validate-certificate

boolean

true

Optional. Set to false to disable certificate validation for the token endpoint call.

solace.oauth.client.auth-server.tls.version

String

TLSv1.2

Optional. The TLS version used for the token endpoint call. Valid values: TLS, TLSv1, TLSv1.1, TLSv1.2, TLSv1.3.

solace.oauth.client.token.refresh.interval

int (seconds)

60

Optional. The number of seconds between token refreshes. Each refresh pushes the token into the live event broker session.

solace.oauth.client.token.fetch.timeout

int (milliseconds)

100

Optional. Connect timeout for the token endpoint HTTP request.

OAuth File Mode

The following configuration options are available for file mode:

Config Option Type Default Value Description

solace.oauth.client.access-token

String (file path)

None

Required for file mode instead of server-fetch mode. Path to a file that another process rotates a valid OAuth token into.

solace.oauth.client.access-token.source

String

file

Optional. Only file is currently supported.

solace.oauth.client.token.refresh.interval

int (seconds)

60

Optional. The number of seconds between file re-reads. Each re-read pushes the token into the live event broker session.

Reading from an Event Broker (Source)

Use spark.readStream.format("solace") to stream messages from an event broker queue into a Spark DataFrame.

The source cannot consume from a partitioned queue.

Config Option Type Valid Values Default Value Description

queue

String

None

Required. The name of the queue to consume from. The queue must already exist on the event broker.

queue.receiveWaitTimeout

int (milliseconds)

10000

Optional. The maximum time to wait for messages before yielding whatever is available to the current micro-batch.

batchSize

int

≥ 0

0

Optional. The maximum number of messages each partition reads per micro-batch. The total across a micro-batch can reach batchSize × partitions. A value of 0 (default) means no per-partition limit. We strongly recommend setting it explicitly. As a starting point, set the queue's Maximum Delivered Unacknowledged Messages per Flow to roughly twice this value; for large message sizes, setting it equal to this value instead can reduce memory usage.

partitions

int

1

Optional. The number of concurrent consumers on the queue. Each partition opens its own flow to the same queue (non-exclusive access), so partitions act as competing consumers. Set to 0 to automatically create one consumer per Spark worker node.

offsetIndicator

String

MESSAGE_ID, CORRELATION_ID, APPLICATION_MESSAGE_ID, or a custom user-property key

MESSAGE_ID

Optional. Which field identifies a message for deduplication and checkpointing. MESSAGE_ID uses the event broker's replication group message ID.

ackLastProcessedMessages

boolean

false

Optional. On restart, acknowledge incoming messages that match the last checkpoint instead of reprocessing them. The connector ignores this option if you set replayStrategy. This is a best-effort deduplication check. The connector has no visibility into the downstream system's state, so we recommend also deduplicating at the downstream system for reliable results.

ignoreCheckpointMessageIdComparisonError

boolean

false

Optional. When the checkpointed message ID cannot be compared with an incoming message's ID, the connector logs an error and processes the message without deduplication, instead of failing the query. This setting has no effect unless you also set ackLastProcessedMessages to true; the connector ignores it if you set replayStrategy.

includeHeaders

boolean

false

Optional. When true, adds a Headers map column to the DataFrame schema. See Message Schema.

closeReceiversOnPartitionClose

boolean

false

Optional. Recreates the consumer flow every time a partition closes.

sparkRuntimePlatform

String

DATABRICKS, OTHER

DATABRICKS

Optional. Controls whether Databricks Unity Catalog Volume checkpoint handling is active. Set to OTHER outside of Databricks.

connectIdleTimeoutInMillis

int

0 (disabled)

Optional. Closes the connection after this many milliseconds of inactivity. Disabled when set to 0.

connectIdleTimeoutCheckInMillis

int

0 (disabled)

Optional. The interval, in milliseconds, between checks for idle-connection timeout. Disabled when set to 0.

Sizing Consumers

Since connector version 3.1.0, connections open from the Spark executor nodes rather than the driver. The connector hashes <queueName>-<index> and asks Spark to schedule each consumer on a preferred executor, so a given consumer tends to land on the same node across batches; this hash-based placement can concentrate consumers on fewer nodes than you might expect. Because of this, we recommend that you size partitions to the number of available executor cores. For example, an executor node with 4 CPU cores supports up to 4 concurrent consumers. The partitions parameter sets the minimum number of consumers; Spark determines the actual maximum based on available task slots, typically observed as roughly max consumers = partitions × total executor nodes.

Setting partitions to 0 lets Spark scale consumers automatically, but this process can create an unbounded number of consumers on the queue and affect event broker connection limits. Avoid this setting unless it's necessary.

With more than one partition, messages are consumed by competing consumers, so the Spark task scheduling, not the order messages arrive on the queue, determines processing order. If your downstream system does not already handle out-of-order data, either preserve order there or use a single consumer (partitions set to 1) per queue and distribute load across multiple queues with an appropriate topic hierarchy instead of increasing partitions on one queue.

Checkpointing

The Solace Connector for Apache Spark tracks read progress in two places: the Spark streaming checkpoint location and a Last Value Queue (LVQ) on the event broker. On restart, the connector can recover its position from the LVQ even if the Spark checkpoint location is unavailable or lost. Even so, point the checkpointLocation option at durable storage rather than local or ephemeral storage, so Spark can also recover reliably from its own checkpoint.

The connector automatically provisions the LVQ if it does not already exist. Depending on who provisions it, different prerequisites apply:

  • If an event broker administrator pre-creates the LVQ, it must be Exclusive, have a Spool Quota of 0, and have a topic subscription already added. The administrator must also explicitly assign the connecting client username as the queue's owner. Creating the queue does not do this automatically. For more information, see Configuring Queue Owners. Set Non-Owner Access to No Access so that no client other than the owner can access the queue; this setting has no effect on the owner itself. The client username's ACL must also permit publish and subscribe access to that topic.

  • If the connector creates the LVQ, the client profile must allow clients to create endpoints, and the ACL must permit publish and subscribe access to the LVQ topic. For an event broker service, enable Allow Client to Create Endpoints in the client profile settings; see Client Profile Settings. For a software event broker or an appliance event broker, enable allow-guaranteed-endpoint-create; see Allowing Clients to Create Guaranteed Endpoints.

Config Option Type Default Value Description

lvq.name

String

solace.spark.connector.state

Optional. The name of the LVQ used to shuttle checkpoint state from worker to driver.

lvq.topic

String

solace/spark/connector/offset

Optional. The topic the LVQ subscribes to.

If you run on Databricks and use a Unity Catalog Volume as the Spark checkpoint location, the connector also needs OAuth machine-to-machine service principal credentials, stored as Databricks secrets:

Config Option Type Default Value Description

databricksSecretScope, databricksHost, databricksClientId, databricksClientSecret

String

None

Required together if the checkpoint location is a Unity Catalog Volume. These are Databricks secret names, not the credential values themselves. The service principal they resolve to must have access to both the target Volume and the secret scope.

databricksClientSecretRefreshInterval

long (minutes)

1440

Optional. The number of minutes between re-reads of the client secret from the secret scope.

Databricks Setup

The connector is a native Apache Spark DataSource V2 implementation and runs on any Spark platform. Nothing in the connector itself is Databricks-specific. The guidance in this section applies if you deploy on Databricks specifically.

The following table lists the Databricks permissions typically needed to install and run the connector.

Task Requirement

Install a Maven library on a cluster

Installing a Maven library requires Can Manage permission on the cluster or a cluster policy that permits Maven libraries.

Install Maven libraries on a Standard or Shared access-mode cluster

A metastore admin must add the coordinate to the workspace's Unity Catalog allowlist first. This is the single most common blocker in a first-time setup; a Dedicated access-mode cluster avoids it entirely.

Read secrets from a secret scope

The connector needs READ permission on the secret scope to read Databricks secrets at runtime.

Use a Unity Catalog Volume as the checkpoint location

The connector needs a workspace-level service principal with access to both the Volume and the secret scope. Executor nodes have no direct access to Volumes, so the connector authenticates to the Volume through Databricks unified authentication using this service principal.

Write to a Delta table target

Writing to a Delta table target requires USE CATALOG, USE SCHEMA, MODIFY, and SELECT permissions on the table.

Cluster Configuration

  • Set the cluster's access mode to Dedicated (single user or group). A Standard or Shared access-mode cluster triggers the Unity Catalog allowlist requirement described in the preceding table.

  • Prefer a Job Compute cluster over an All-Purpose cluster. A Job Compute cluster is ephemeral, and its Spark session closes gracefully on termination, cleanly closing event broker connections. An All-Purpose cluster stays active until manually stopped or idled out, which can leave hanging consumer connections on the event broker and affect connection limits. If you use an All-Purpose cluster, set connectIdleTimeoutInMillis and connectIdleTimeoutCheckInMillis to automatically close inactive connections; this is also useful for intermittent data flows, because consumers then exist only while there's data in the queue. The event broker redelivers messages still pending acknowledgment when a connection closes this way, so ensure your downstream processing is idempotent. This close is best-effort and may not occur if the underlying JVM has trouble during termination.

  • Prefer on-demand instances over spot instances. Abrupt spot-instance termination triggers cluster reinitialization, which can leave old consumers ungracefully terminated on the event broker while new ones are created, affecting connection limits and disrupting data flow.

  • Prefer several smaller job clusters, each handling a specific set of queues grouped by similar characteristics (payload size, ingress rate, expected processing rate), over one large cluster handling everything. A cluster failure then affects only a subset of your workflows. Give queues with large payloads their own dedicated cluster. If a cluster autoscales, use the Databricks-recommended default of a minimum of one worker and a maximum equal to your configured worker count.

  • Before installing a new connector version, remove all earlier versions from the cluster and restart it. A leftover JAR file from a previous version causes classpath conflicts, and this is especially important if you move across the 3.x/4.x boundary, because a stale Scala 2.12 build is binary-incompatible with Scala 2.13 and causes errors such as NoSuchMethodError at startup.

Managing Secrets

Store credentials in a Databricks secret scope using the Databricks CLI, and read them with dbutils.secrets.get. For non-sensitive settings or secret references, you can also use the Environment variables and Spark config fields, which are on the Spark tab under Advanced in the cluster configuration; the Spark config field is also useful for cluster-wide properties, such as spark.sql.ansi.enabled if you migrate to the 4.x line. Restart the cluster after changing values in either field.

Values set in the Environment variables or Spark config fields are visible in plaintext to anyone who can view the cluster configuration. Use them only for non-sensitive settings or secret references, never for secret values themselves.

Secret rotation is not managed by the connector. A platform administrator rotates the secret and updates it in the scope; the connector re-reads it at the interval, in minutes, that you set with databricksClientSecretRefreshInterval, with no restart required. Both the rotation and the connector's next refresh must complete before the old secret expires. If the connector's refresh fires before rotation completes, it keeps using the old secret until the next refresh cycle.

Message Replay

The Solace Connector for Apache Spark can start consuming from a replayed message log instead of the queue's current head, provided the event broker's Replay Log feature is enabled and configured. When you set replayStrategy, the connector ignores ackLastProcessedMessages.

Config Option Type Valid Values Description

replayStrategy

String

BEGINNING, TIMEBASED, REPLICATION-GROUP-MESSAGE-ID

Optional. Not set by default (no replay).

replayReplicationGroupMessageId

String

Required if replayStrategy is REPLICATION-GROUP-MESSAGE-ID. The replication group message ID to replay from.

replayStartTime

String (yyyy-MM-dd'T'HH:mm:ss)

Required if replayStrategy is TIMEBASED. The replay start timestamp.

replayStartTimeTimezone

String

UTC

Optional. The timezone used to interpret replayStartTime.

Writing to an Event Broker (Sink)

Use writeStream (recommended) or forEachBatch to publish DataFrame rows to the event broker. A row missing a Payload value, missing an Id value, or missing both a Topic value and the topic option raises an exception rather than publishing with a default.

Every published message uses PERSISTENT delivery.

Prefer writeStream over forEachBatch for authentication schemes other than basic username/password. The batch-write path currently requires you to set username and password even when you configure client certificate or OAuth authentication.

Config Option Type Valid Values Default Value Description

topic

String

None

Optional, but required if the DataFrame has no Topic column. A static topic that the connector publishes every row to. See Dynamic Destinations.

includeHeaders

boolean

false

Optional. When true, publishes the row's Headers map column as message user properties.

publishAckTimeout

long (milliseconds)

5000

Optional. How long to wait, per commit, for all publish acknowledgments to arrive. Tune against network latency to the event broker, the number of rows in each micro-batch or forEachBatch invocation, and expected event broker response time under load.

publishAckTimeoutFailOnError

boolean

true

Optional. When true, the connector throws an exception on an acknowledgment timeout. When false, the connector logs the timeout and continues processing.

Dynamic Destinations

If you don't set the topic option, the connector publishes each row to the topic in its Topic column, letting you route different rows to different topics from the same write.

Message Schema

The source produces DataFrames with this fixed schema; it is not extensible beyond the includeHeaders toggle. The sink does not require a matching schema: it reads the following columns from your DataFrame by name (case-sensitive) and ignores any others.

Column Spark Type Description

Id

StringType

On read, the value identified by offsetIndicator. On write, the application message ID; required for every row.

Payload

BinaryType

The message body, as raw bytes. Only the byte[] type is accepted on write. Other payload types raise an error. The connector enforces no size limit of its own; the underlying event broker message-size limit applies.

Topic

StringType

On read, the topic on which the connector received the message. On write, the destination topic if the topic option is not set.

TimeStamp

TimestampType

The sender timestamp, or the receive timestamp if the publisher does not set the sender timestamp.

Headers

MapType(StringType, BinaryType)

Present only when includeHeaders is true. On read, includes JCSMP message user properties plus connector-added entries (sequence number, expiration, time-to-live, priority, redelivered flag). On write, the connector maps keys matching a known set of message fields that use the solace_ prefix (for example solace_applicationMessageId, solace_correlationId, solace_priority, solace_timeToLive) to those fields with automatic type coercion; other keys become generic user properties.

If you don't need to replay from a specific message later, drop the Id column before writing to your output sink. In benchmarking, dropping this column alone cut Spark per-batch write (addBatch) time from about 3 seconds to 1-2 seconds and raised throughput by roughly 1.5×. This applies whether replay is connector-initiated or initiated by an event broker administrator. Neither needs the Id column retained downstream.

Performance Considerations

Several factors affect the throughput you can expect from the connector:

  • batchSize and consumer count together

  • The queue's Maximum Delivered Unacknowledged Messages per Flow relative to batchSize (set it to roughly twice batchSize for stable throughput)

  • Whether the Id column is dropped before the output sink

  • Message size

  • Network latency to the event broker

  • The write speed of your output sink

  • Any downstream transformations

  • On Databricks, whether Photon is enabled