Cloud

gcp_pubsub

Sends messages to a GCP Cloud Pub/Sub topic. Metadata from messages are sent as attributes.

For information on how to set up credentials, see this guide.

Troubleshooting

If you’re consistently seeing Failed to send message to gcp_pubsub: context deadline exceeded error logs without any further information it is possible that you are encountering https://github.com/redpanda-data/connect/issues/1042, which occurs when metadata values contain characters that are not valid utf-8. This can frequently occur when consuming from Kafka as the key metadata field may be populated with an arbitrary binary value, but this issue is not exclusive to Kafka.

If you are blocked by this issue then a work around is to delete either the specific problematic keys:

pipeline:
  processors:
    - mapping: |
        meta kafka_key = deleted()

Or delete all keys with:

pipeline:
  processors:
    - mapping: meta = deleted()
  • Common

  • Advanced

output:
  label: ""
  gcp_pubsub:
    project: "" # No default (required)
    credentials_json: ""
    topic: "" # No default (required)
    endpoint: ""
    max_in_flight: 64
    count_threshold: 100
    delay_threshold: 10ms
    byte_threshold: 1000000
    metadata:
      exclude_prefixes: []
    batching:
      count: 0
      byte_size: 0
      period: ""
      check: ""
output:
  label: ""
  gcp_pubsub:
    project: "" # No default (required)
    credentials_json: ""
    topic: "" # No default (required)
    endpoint: ""
    ordering_key: "" # No default (optional)
    max_in_flight: 64
    count_threshold: 100
    delay_threshold: 10ms
    byte_threshold: 1000000
    publish_timeout: 1m0s
    validate_topic: true
    metadata:
      exclude_prefixes: []
    flow_control:
      max_outstanding_bytes: -1
      max_outstanding_messages: 1000
      limit_exceeded_behavior: block
    batching:
      count: 0
      byte_size: 0
      period: ""
      check: ""
      processors: [] # No default (optional)

Fields

batching

Configures a batching policy on this output. While the PubSub client maintains its own internal buffering mechanism, preparing larger batches of messages can further trade-off some latency for throughput.

Type: object

# Examples:
batching:
  byte_size: 5000
  count: 0
  period: 1s

# ---

batching:
  count: 10
  period: 1s

# ---

batching:
  check: this.contains("END BATCH")
  count: 0
  period: 1m

batching.byte_size

The maximum total size (in bytes) that a batch can reach before it is flushed. When the combined size of all messages in the batch reaches or exceeds this limit, the batch is immediately sent to the next stage (such as a processor or output).

Set to 0 to disable size-based batching. When disabled, messages are flushed based on other conditions (such as count or period).

Type: int

Default: 0

batching.check

A Bloblang query that returns a boolean value indicating whether a message should end a batch.

Type: string

Default: ""

# Examples:
check: this.type == "end_of_transaction"

batching.count

The number of messages at which the batch is flushed. Set to 0 to disable count-based batching.

Type: int

Default: 0

batching.period

The length of time after which an incomplete batch is flushed regardless of its size. This field accepts Go duration format strings such as 100ms, 1s, or 5s. Supported time units are ns, us, ms, s, m, and h.

Type: string

Default: ""

# Examples:
period: 1s

# ---

period: 1m

# ---

period: 500ms

batching.processors[]

A list of processors to apply to a batch as it is flushed. This allows you to aggregate and archive the batch however you see fit. All resulting messages are flushed as a single batch, so splitting the batch into smaller batches with these processors has no effect.

Type: array<processor>

# Examples:
processors:
  - archive:
      format: concatenate


# ---

processors:
  - archive:
      format: lines


# ---

processors:
  - archive:
      format: json_array

byte_threshold

Publish a batch when its size in bytes reaches this value.

Type: int

Default: 1000000

count_threshold

Publish a pubsub buffer when it has this many messages

Type: int

Default: 100

credentials_json

The Google Service Account credentials in JSON format (optional). Provide the contents of the credentials file as plain JSON, not Base64-encoded. Use this field to authenticate with Google Cloud services. If this field is empty, the component uses Application Default Credentials. For more information about creating service account credentials, see Google’s service account documentation.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

delay_threshold

Publish a non-empty pubsub buffer after this delay has passed.

Type: string

Default: 10ms

endpoint

An optional endpoint that overrides the default of pubsub.googleapis.com:443. Use this field to connect to a region-specific Pub/Sub endpoint. For a list of valid values, see Pub/Sub regional endpoints.

Type: string

Default: ""

# Examples:
endpoint: us-central1-pubsub.googleapis.com:443

# ---

endpoint: us-west3-pubsub.googleapis.com:443

flow_control

For a given topic, configures the PubSub client’s internal buffer for messages to be published.

Type: object

flow_control.limit_exceeded_behavior

Configures the behavior when trying to publish additional messages while the flow controller is full. The available options are block (default), ignore (disable), and signal_error (publish results will return an error).

Type: string

Default: block

Options: ignore, block, signal_error

flow_control.max_outstanding_bytes

Maximum size of buffered messages to be published. If less than or equal to zero, this is disabled.

Type: int

Default: -1

flow_control.max_outstanding_messages

Maximum number of buffered messages to be published. If less than or equal to zero, this is disabled.

Type: int

Default: 1000

max_in_flight

The maximum number of messages to have in flight at a given time. Increasing this may improve throughput.

Type: int

Default: 64

metadata

Specify criteria for which metadata values are sent as attributes, all are sent by default.

Type: object

metadata.exclude_prefixes[]

Provide a list of explicit metadata key prefixes to exclude when adding metadata to sent messages.

Type: array<string>

Default: []

ordering_key

The ordering key to use for publishing messages.

This field supports interpolation functions.

Type: string

project

The project ID of the topic to publish to.

Type: string

publish_timeout

The maximum length of time to wait before abandoning a publish attempt for a message.

Type: string

Default: 1m0s

# Examples:
publish_timeout: 10s

# ---

publish_timeout: 5m

# ---

publish_timeout: 60m

topic

The topic to publish to.

This field supports interpolation functions.

Type: string

validate_topic

Whether to validate the existence of the topic before publishing. If set to false and the topic does not exist, messages will be lost.

Type: bool

Default: true