Managed I/O Connectors

Beam’s new Managed API streamlines how you use existing I/Os, offering both simplicity and powerful enhancements. I/Os are now configured through a lightweight, consistent interface: a simple configuration map with a unified API that spans multiple connectors.

With Managed I/O, runners gain deeper insight into each I/O’s structure and intent. This allows the runner to optimize performance, adjust behavior dynamically, or even replace the I/O with a more efficient or updated implementation behind the scenes.

For example, the DataflowRunner can seamlessly upgrade a Managed transform to its latest SDK version, automatically applying bug fixes and new features (no manual updates or user intervention required!)

Supported SDKs

The Managed API is directly accessible through the Java and Python SDKs.

Additionally, some SDKs use the Managed API internally. For example, the Iceberg connector used in Beam YAML and Beam SQL is invoked via the Managed API under the hood.

Available Configurations

Note: required configuration fields are bolded.

Connector NameRead ConfigurationWrite Configuration
ICEBERGtable (str)
catalog_name (str)
catalog_properties (map[str, str])
config_properties (map[str, str])
drop (list[str])
filter (str)
keep (list[str])
table (str)
allowed_lateness_seconds (int32)
autosharding (boolean)
catalog_name (str)
catalog_properties (map[str, str])
change_type_column (str)
change_type_map (map[str, str])
config_properties (map[str, str])
direct_write_byte_limit (int32)
distribution_mode (str)
drop (list[str])
equality_columns (list[str])
keep (list[str])
maximum_table_cache_size (int32)
mode (str)
num_shards (int32)
only (str)
partition_fields (list[str])
sequence_number_column (str)
shards_per_partition (int32)
sink_id (str)
snapshot_properties (map[str, str])
sort_fields (list[str])
sorter_memory_mb (int32)
table_cache_polling_buckets (int32)
table_cache_refresh_interval_seconds (int32)
table_properties (map[str, str])
token_heartbeat_seconds (int32)
triggering_frequency_seconds (int32)
upsert (boolean)
use_side_input_table_cache (boolean)
write_properties (map[str, str])
DELTA_CDCtable (str)
end_timestamp (str)
end_version (int64)
hadoop_config (map[str, str])
include_metadata_columns (list[str])
start_timestamp (str)
start_version (int64)
Unavailable
DELTAtable (str)
hadoop_config (map[str, str])
timestamp (str)
version (int64)
Unavailable
KAFKAbootstrap_servers (str)
topic (str)
allow_duplicates (boolean)
confluent_schema_registry_subject (str)
confluent_schema_registry_url (str)
consumer_config_updates (map[str, str])
file_descriptor_path (str)
format (str)
message_name (str)
offset_deduplication (boolean)
redistribute_by_record_key (boolean)
redistribute_num_keys (int32)
redistributed (boolean)
schema (str)
bootstrap_servers (str)
format (str)
topic (str)
file_descriptor_path (str)
message_name (str)
producer_config_updates (map[str, str])
schema (str)
ICEBERG_CDCtable (str)
catalog_name (str)
catalog_properties (map[str, str])
config_properties (map[str, str])
drop (list[str])
filter (str)
from_snapshot (int64)
from_timestamp (int64)
include_metadata_columns (list[str])
keep (list[str])
poll_interval_seconds (int32)
starting_strategy (str)
streaming (boolean)
to_snapshot (int64)
to_timestamp (int64)
watermark_column (str)
watermark_column_time_unit (str)
Unavailable
MYSQLjdbc_url (str)
connection_init_sql (list[str])
connection_properties (str)
disable_auto_commit (boolean)
fetch_size (int32)
location (str)
num_partitions (int32)
output_parallelization (boolean)
partition_column (str)
password (str)
read_query (str)
secret_manager (str)
username (str)
jdbc_url (str)
autosharding (boolean)
batch_size (int64)
connection_init_sql (list[str])
connection_properties (str)
location (str)
password (str)
secret_manager (str)
username (str)
write_statement (str)
POSTGRESjdbc_url (str)
connection_properties (str)
fetch_size (int32)
location (str)
num_partitions (int32)
output_parallelization (boolean)
partition_column (str)
password (str)
read_query (str)
secret_manager (str)
username (str)
jdbc_url (str)
autosharding (boolean)
batch_size (int64)
connection_properties (str)
location (str)
password (str)
secret_manager (str)
username (str)
write_statement (str)
BIGQUERYkms_key (str)
query (str)
row_restriction (str)
fields (list[str])
table (str)
table (str)
drop (list[str])
keep (list[str])
kms_key (str)
only (str)
triggering_frequency_seconds (int64)
SQLSERVERjdbc_url (str)
connection_properties (str)
disable_auto_commit (boolean)
fetch_size (int32)
location (str)
num_partitions (int32)
output_parallelization (boolean)
partition_column (str)
password (str)
read_query (str)
secret_manager (str)
username (str)
jdbc_url (str)
autosharding (boolean)
batch_size (int64)
connection_properties (str)
location (str)
password (str)
secret_manager (str)
username (str)
write_statement (str)

Configuration Details

ICEBERG Write

’).

ConfigurationTypeDescription
tablestrA fully-qualified table identifier. You may also provide a template to write to multiple dynamic destinations, for example: `dataset.my_{col1}_{col2.nested}_table`.
allowed_lateness_secondsint32How long a late record may lag behind the watermark before it is dropped entirely, rather than routed to the dead_letter output. Defaults to 21600 (6 hours). Currently only supported in 'merge-on-read' mode.
autoshardingbooleanEnables dynamic sharding to automatically adjust the number of parallel writers based on data volume. It handles data skew by further sub-dividing partitions into multiple shards to prevent bottlenecks during high-throughput writes. Only available with 'hash' distribution mode.
catalog_namestrName of the catalog containing the table.
catalog_propertiesmap[str, str]Properties used to set up the Iceberg catalog.
change_type_columnstrMerge-on-read only. The optional column name representing the row's change type (INSERT, UPDATE_BEFORE, UPDATE_AFTER, or DELETE). This column will be stripped from the data row before writing to Iceberg. If unset, the sink will use the element's native ValueKind
change_type_mapmap[str, str]Merge-on-read only. Optional map from a change_type_column value to the canonical change type name (see above).
config_propertiesmap[str, str]Properties passed to the Hadoop Configuration.
direct_write_byte_limitint32For a streaming pipeline, sets the limit for lifting bundles into the direct write path.
distribution_modestrDefines distribution of write data. Supported distributions: - none: don't shuffle rows (default) - hash: shuffle rows by partition key before writing data
droplist[str]A list of field names to drop from the input record before writing. Is mutually exclusive with 'keep' and 'only'. In merge-on-read mode the control columns are always dropped.
equality_columnslist[str]Columns defining row identity (equality-delete fields). Defaults to the destination table's identifier (primary-key) fields. Required if the table doesn't exist yet. Currently only supported in 'merge-on-read' mode.
keeplist[str]A list of field names to keep in the input record. All other fields are dropped before writing. Is mutually exclusive with 'drop' and 'only'. In merge-on-read mode the control columns are dropped unless listed here.
maximum_table_cache_sizeint32For a batch pipeline, sets the maximum number of table metadata specs to cache in memory. Tables exceeding this limit fall back to worker-local catalog loading.
modestrControls how rows are written. 'append' (default) appends every row as new data. 'merge-on-read' treats each row as a change (INSERT, UPDATE_BEFORE, UPDATE_AFTER, or DELETE) applied to the table by primary key.
num_shardsint32The number of deterministic primary-key-hash shards per destination, i.e. the max write parallelism per destination. Too low may bottleneck writes, and too high may produce more files. Defaults to 16. Currently only supported in 'merge-on-read' mode.
onlystrThe name of a single record field that should be written. Is mutually exclusive with 'keep' and 'drop'.
partition_fieldslist[str]Fields used to create a partition spec that is applied when tables are created. For a field 'foo', the available partition transforms are:
  • foo
  • truncate(foo, N)
  • bucket(foo, N)
  • hour(foo)
  • day(foo)
  • month(foo)
  • year(foo)
  • void(foo)

For more information on partition transforms, please visit https://iceberg.apache.org/spec/#partition-transforms.

sequence_number_columnstrMerge-on-read only. The required column name representing the monotonic sequence number used to order a single key’s changes. Defaults to ‘_commit_snapshot_sequence_number’. This column will be stripped from the data row before writing to Iceberg.
shards_per_partitionint32Maximum number of shards a single partition’s rows may occupy. Lower values write fewer files per commit, but also reduces per-partition write parallelism. A value of 1 pins each partition to one writer. Ignored for unpartitioned tables. Must be between 1 and num_shards; defaults to num_shards. Currently only supported in ‘merge-on-read’ mode.
sink_idstrA stable identifier for this sink, used to namespace the idempotency tokens written to each commit’s Iceberg snapshot summary. Defaults to a unique per-write UUID. Set it explicitly (and keep it stable across relaunches) for exactly-once commits across relaunches of a particular streaming write. A batch load with a stable sink_id commits only once (later batch loads with the same sink_id are skipped). Currently only supported in ‘merge-on-read’ mode.
snapshot_propertiesmap[str, str]Extra key/value properties to add to every commit’s Iceberg snapshot summary. Keys prefixed with ‘beam.cdc.’ are reserved and rejected. Currently only supported in ‘merge-on-read’ mode.
sort_fieldslist[str]Fields used to set the table’s sort order, applied when the table is created. Each entry has the form <term> [asc|desc] [nulls first|nulls last], where <term> is a field name or one of the partition transforms (e.g. bucket(col, 4), day(ts)). Direction defaults to ascending; null order defaults to nulls-first for ascending and nulls-last for descending. Note: this sets the table’s declared sort order as metadata; it does not cause Beam to physically sort records before writing. For more information on sort orders, please visit https://iceberg.apache.org/spec/#sort-orders.
sorter_memory_mbint32The in-memory buffer size (MB) for the pre-write sort; groups larger than this spill to disk. Must be >= 1. Defaults to 100. Currently only supported in ‘merge-on-read’ mode.
table_cache_polling_bucketsint32Sets the number of parallel buckets/workers used to query the Iceberg catalog during refreshes. Defaults to 1.
table_cache_refresh_interval_secondsint32For a streaming pipeline, sets the interval in seconds at which table metadata is refreshed from the catalog.
table_propertiesmap[str, str]Iceberg table properties to be set on the table when it is created. For more information on table properties, please visit https://iceberg.apache.org/docs/latest/configuration/#table-properties.
token_heartbeat_secondsint32Streaming only. If set, the sink will emit a periodic empty token-refresh commit while idle, so its thread of sink_id stamped snapshot stays recent and is less likely to be lost to expire_snapshots. Disabled by default. Currently only supported in ‘merge-on-read’ mode.
triggering_frequency_secondsint32For a streaming pipeline, sets the frequency at which snapshots are produced.
upsertbooleanMerge-on-read only. If true, only the after-image of each change (INSERT/UPDATE_AFTER) is applied, as an upsert; UPDATE_BEFORE records are dropped. Default: false.
use_side_input_table_cachebooleanEnables expirable side-input caching of Iceberg table metadata across workers to reduce catalog load.
write_propertiesmap[str, str]Properties applied to the underlying file writer (e.g. Parquet write properties like ‘write.parquet.bloom-filter-enabled.column.

ICEBERG Read

ConfigurationTypeDescription
tablestrIdentifier of the Iceberg table.
catalog_namestrName of the catalog containing the table.
catalog_propertiesmap[str, str]Properties used to set up the Iceberg catalog.
config_propertiesmap[str, str]Properties passed to the Hadoop Configuration.
droplist[str]A subset of column names to exclude from reading. If null or empty, all columns will be read.
filterstrSQL-like predicate to filter data at scan time. Example: "id > 5 AND status = 'ACTIVE'". Uses Apache Calcite syntax: https://calcite.apache.org/docs/reference.html
keeplist[str]A subset of column names to read exclusively. If null or empty, all columns will be read.

DELTA_CDC Read

ConfigurationTypeDescription
tablestrIdentifier of the Delta Lake table.
end_timestampstrEnd timestamp of the Delta Lake table to read changes up to. Should be specified in the ISO 8601 standard.
end_versionint64End version of the Delta Lake table to read changes up to.
hadoop_configmap[str, str]Properties passed to the Hadoop Configuration.
include_metadata_columnslist[str]Metadata columns to include in the output rows. Supported columns are: _change_type, _commit_version, and _commit_timestamp.
start_timestampstrStart timestamp of the Delta Lake table to read changes from. Should be specified in the ISO 8601 standard. Either this or the start version has to be provided.
start_versionint64Start version of the Delta Lake table to read changes from. Either this or the start timestamp has to be provided.

DELTA Read

ConfigurationTypeDescription
tablestrIdentifier of the Delta Lake table.
hadoop_configmap[str, str]Properties passed to the Hadoop Configuration.
timestampstrTimestamp of the Delta Lake table to read (in UTC ISO 8601 format, e.g. 2026-05-20T15:43:26Z). Cannot be set if version is set.
versionint64Version of the Delta Lake table to read. Cannot be set if timestamp is set.

KAFKA Read

ConfigurationTypeDescription
bootstrap_serversstrA list of host/port pairs to use for establishing the initial connection to the Kafka cluster. The client will make use of all servers irrespective of which servers are specified here for bootstrapping—this list only impacts the initial hosts used to discover the full set of servers. This list should be in the form `host1:port1,host2:port2,...`
topicstrn/a
allow_duplicatesbooleanIf the Kafka read allows duplicates.
confluent_schema_registry_subjectstrn/a
confluent_schema_registry_urlstrn/a
consumer_config_updatesmap[str, str]A list of key-value pairs that act as configuration parameters for Kafka consumers. Most of these configurations will not be needed, but if you need to customize your Kafka consumer, you may use this. See a detailed list: https://docs.confluent.io/platform/current/installation/configuration/consumer-configs.html
file_descriptor_pathstrThe path to the Protocol Buffer File Descriptor Set file. This file is used for schema definition and message serialization.
formatstrThe encoding format for the data stored in Kafka. Valid options are: RAW,STRING,AVRO,JSON,PROTO
message_namestrThe name of the Protocol Buffer message to be used for schema extraction and data conversion.
offset_deduplicationbooleanIf the redistribute is using offset deduplication mode.
redistribute_by_record_keybooleanIf the redistribute keys by the Kafka record key.
redistribute_num_keysint32The number of keys for redistributing Kafka inputs.
redistributedbooleanIf the Kafka read should be redistributed.
schemastrThe schema in which the data is encoded in the Kafka topic. For AVRO data, this is a schema defined with AVRO schema syntax (https://avro.apache.org/docs/1.10.2/spec.html#schemas). For JSON data, this is a schema defined with JSON-schema syntax (https://json-schema.org/). If a URL to Confluent Schema Registry is provided, then this field is ignored, and the schema is fetched from Confluent Schema Registry.

KAFKA Write

ConfigurationTypeDescription
bootstrap_serversstrA list of host/port pairs to use for establishing the initial connection to the Kafka cluster. The client will make use of all servers irrespective of which servers are specified here for bootstrapping—this list only impacts the initial hosts used to discover the full set of servers. | Format: host1:port1,host2:port2,...
formatstrThe encoding format for the data stored in Kafka. Valid options are: RAW,JSON,AVRO,PROTO
topicstrn/a
file_descriptor_pathstrThe path to the Protocol Buffer File Descriptor Set file. This file is used for schema definition and message serialization.
message_namestrThe name of the Protocol Buffer message to be used for schema extraction and data conversion.
producer_config_updatesmap[str, str]A list of key-value pairs that act as configuration parameters for Kafka producers. Most of these configurations will not be needed, but if you need to customize your Kafka producer, you may use this. See a detailed list: https://docs.confluent.io/platform/current/installation/configuration/producer-configs.html
schemastrn/a

ICEBERG_CDC Read

ConfigurationTypeDescription
tablestrIdentifier of the Iceberg table.
catalog_namestrName of the catalog containing the table.
catalog_propertiesmap[str, str]Properties used to set up the Iceberg catalog.
config_propertiesmap[str, str]Properties passed to the Hadoop Configuration.
droplist[str]A subset of column names to exclude from reading. If null or empty, all columns will be read.
filterstrSQL-like predicate to filter data at scan time. Example: "id > 5 AND status = 'ACTIVE'". Uses Apache Calcite syntax: https://calcite.apache.org/docs/reference.html
from_snapshotint64Starts reading from this snapshot ID (inclusive).
from_timestampint64Starts reading from the first snapshot (inclusive) that was created after this timestamp (in milliseconds).
include_metadata_columnslist[str]List of top-level metadata columns to include with CDC output rows. Supported columns: - `_change_type` - `_row_id` - `_last_updated_sequence_number` - `_commit_snapshot_id` - `_commit_snapshot_sequence_number`
  </td>
</tr>
<tr>
  <td>
    keep
  </td>
  <td>
    <code>list[<span style="color: green;">str</span>]</code>
  </td>
  <td>
    A subset of column names to read exclusively. If null or empty, all columns will be read.
  </td>
</tr>
<tr>
  <td>
    poll_interval_seconds
  </td>
  <td>
    <code style="color: #f54251">int32</code>
  </td>
  <td>
    The interval at which to poll for new snapshots. Defaults to 60 seconds.
  </td>
</tr>
<tr>
  <td>
    starting_strategy
  </td>
  <td>
    <code style="color: green">str</code>
  </td>
  <td>
    The source's starting strategy. Valid options are: "earliest" or "latest". Can be overriden by setting a starting snapshot or timestamp. Defaults to earliest for batch, and latest for streaming.
  </td>
</tr>
<tr>
  <td>
    streaming
  </td>
  <td>
    <code style="color: orange">boolean</code>
  </td>
  <td>
    Enables streaming reads, where source continuously polls for snapshots forever.
  </td>
</tr>
<tr>
  <td>
    to_snapshot
  </td>
  <td>
    <code style="color: #f54251">int64</code>
  </td>
  <td>
    Reads up to this snapshot ID (inclusive).
  </td>
</tr>
<tr>
  <td>
    to_timestamp
  </td>
  <td>
    <code style="color: #f54251">int64</code>
  </td>
  <td>
    Reads up to the latest snapshot (inclusive) created before this timestamp (in milliseconds).
  </td>
</tr>
<tr>
  <td>
    watermark_column
  </td>
  <td>
    <code style="color: green">str</code>
  </td>
  <td>
    Column used to derive the source's output watermark. Must be an existing, required, top-level column of type 'long' or 'timestamp'. If not set, the watermark advances according to snapshot commit timestamp.
  </td>
</tr>
<tr>
  <td>
    watermark_column_time_unit
  </td>
  <td>
    <code style="color: green">str</code>
  </td>
  <td>
    Time unit used to interpret watermark column of type LONG. One of NANOSECONDS, MICROSECONDS, MILLISECONDS, SECONDS, MINUTES, HOURS, DAYS. Defaults to MICROSECONDS.
  </td>
</tr>

MYSQL Write

ConfigurationTypeDescription
jdbc_urlstrConnection URL for the JDBC sink.
autoshardingbooleanIf true, enables using a dynamically determined number of shards to write.
batch_sizeint64n/a
connection_init_sqllist[str]Sets the connection init sql statements used by the Driver. Only MySQL and MariaDB support this.
connection_propertiesstrUsed to set connection properties passed to the JDBC driver not already defined as standalone parameter (e.g. username and password can be set using parameters above accordingly). Format of the string must be "key1=value1;key2=value2;".
locationstrName of the table to write to.
passwordstrPassword for the JDBC source. Can be specified as a plain password, or as a secret specification in JSON format if used with a secret manager.
secret_managerstrSecret Manager to use for fetching secret values. Available options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'. If not set, no secret manager is used and the password is treated as a plain password.
usernamestrUsername for the JDBC source.
write_statementstrSQL query used to insert records into the JDBC sink.

MYSQL Read

ConfigurationTypeDescription
jdbc_urlstrConnection URL for the JDBC source.
connection_init_sqllist[str]Sets the connection init sql statements used by the Driver. Only MySQL and MariaDB support this.
connection_propertiesstrUsed to set connection properties passed to the JDBC driver not already defined as standalone parameter (e.g. username and password can be set using parameters above accordingly). Format of the string must be "key1=value1;key2=value2;".
disable_auto_commitbooleanWhether to disable auto commit on read. Defaults to true if not provided. The need for this config varies depending on the database platform. Informix requires this to be set to false while Postgres requires this to be set to true.
fetch_sizeint32This method is used to override the size of the data that is going to be fetched and loaded in memory per every database call. It should ONLY be used if the default value throws memory errors.
locationstrName of the table to read from.
num_partitionsint32The number of partitions
output_parallelizationbooleanWhether to reshuffle the resulting PCollection so results are distributed to all workers.
partition_columnstrName of a column of numeric type that will be used for partitioning.
passwordstrPassword for the JDBC source. Can be specified as a plain password, or as a secret specification in JSON format if used with a secret manager.
read_querystrSQL query used to query the JDBC source.
secret_managerstrSecret Manager to use for fetching secret values. Available options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'. If not set, no secret manager is used and the password is treated as a plain password.
usernamestrUsername for the JDBC source.

POSTGRES Write

ConfigurationTypeDescription
jdbc_urlstrConnection URL for the JDBC sink.
autoshardingbooleanIf true, enables using a dynamically determined number of shards to write.
batch_sizeint64n/a
connection_propertiesstrUsed to set connection properties passed to the JDBC driver not already defined as standalone parameter (e.g. username and password can be set using parameters above accordingly). Format of the string must be "key1=value1;key2=value2;".
locationstrName of the table to write to.
passwordstrPassword for the JDBC source. Can be specified as a plain password, or as a secret specification in JSON format if used with a secret manager.
secret_managerstrSecret Manager to use for fetching secret values. Available options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'. If not set, no secret manager is used and the password is treated as a plain password.
usernamestrUsername for the JDBC source.
write_statementstrSQL query used to insert records into the JDBC sink.

POSTGRES Read

ConfigurationTypeDescription
jdbc_urlstrConnection URL for the JDBC source.
connection_propertiesstrUsed to set connection properties passed to the JDBC driver not already defined as standalone parameter (e.g. username and password can be set using parameters above accordingly). Format of the string must be "key1=value1;key2=value2;".
fetch_sizeint32This method is used to override the size of the data that is going to be fetched and loaded in memory per every database call. It should ONLY be used if the default value throws memory errors.
locationstrName of the table to read from.
num_partitionsint32The number of partitions
output_parallelizationbooleanWhether to reshuffle the resulting PCollection so results are distributed to all workers.
partition_columnstrName of a column of numeric type that will be used for partitioning.
passwordstrPassword for the JDBC source. Can be specified as a plain password, or as a secret specification in JSON format if used with a secret manager.
read_querystrSQL query used to query the JDBC source.
secret_managerstrSecret Manager to use for fetching secret values. Available options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'. If not set, no secret manager is used and the password is treated as a plain password.
usernamestrUsername for the JDBC source.

BIGQUERY Read

ConfigurationTypeDescription
kms_keystrUse this Cloud KMS key to encrypt your data
querystrThe SQL query to be executed to read from the BigQuery table.
row_restrictionstrRead only rows that match this filter, which must be compatible with Google standard SQL. This is not supported when reading via query.
fieldslist[str]Read only the specified fields (columns) from a BigQuery table. Fields may not be returned in the order specified. If no value is specified, then all fields are returned. Example: "col1, col2, col3"
tablestrThe fully-qualified name of the BigQuery table to read from. Format: [${PROJECT}:]${DATASET}.${TABLE}

BIGQUERY Write

ConfigurationTypeDescription
tablestrThe bigquery table to write to. Format: [${PROJECT}:]${DATASET}.${TABLE}
droplist[str]A list of field names to drop from the input record before writing. Is mutually exclusive with 'keep' and 'only'.
keeplist[str]A list of field names to keep in the input record. All other fields are dropped before writing. Is mutually exclusive with 'drop' and 'only'.
kms_keystrUse this Cloud KMS key to encrypt your data
onlystrThe name of a single record field that should be written. Is mutually exclusive with 'keep' and 'drop'.
triggering_frequency_secondsint64Determines how often to 'commit' progress into BigQuery. Default is every 5 seconds.

SQLSERVER Read

ConfigurationTypeDescription
jdbc_urlstrConnection URL for the JDBC source.
connection_propertiesstrUsed to set connection properties passed to the JDBC driver not already defined as standalone parameter (e.g. username and password can be set using parameters above accordingly). Format of the string must be "key1=value1;key2=value2;".
disable_auto_commitbooleanWhether to disable auto commit on read. Defaults to true if not provided. The need for this config varies depending on the database platform. Informix requires this to be set to false while Postgres requires this to be set to true.
fetch_sizeint32This method is used to override the size of the data that is going to be fetched and loaded in memory per every database call. It should ONLY be used if the default value throws memory errors.
locationstrName of the table to read from.
num_partitionsint32The number of partitions
output_parallelizationbooleanWhether to reshuffle the resulting PCollection so results are distributed to all workers.
partition_columnstrName of a column of numeric type that will be used for partitioning.
passwordstrPassword for the JDBC source. Can be specified as a plain password, or as a secret specification in JSON format if used with a secret manager.
read_querystrSQL query used to query the JDBC source.
secret_managerstrSecret Manager to use for fetching secret values. Available options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'. If not set, no secret manager is used and the password is treated as a plain password.
usernamestrUsername for the JDBC source.

SQLSERVER Write

ConfigurationTypeDescription
jdbc_urlstrConnection URL for the JDBC sink.
autoshardingbooleanIf true, enables using a dynamically determined number of shards to write.
batch_sizeint64n/a
connection_propertiesstrUsed to set connection properties passed to the JDBC driver not already defined as standalone parameter (e.g. username and password can be set using parameters above accordingly). Format of the string must be "key1=value1;key2=value2;".
locationstrName of the table to write to.
passwordstrPassword for the JDBC source. Can be specified as a plain password, or as a secret specification in JSON format if used with a secret manager.
secret_managerstrSecret Manager to use for fetching secret values. Available options: 'GoogleCloudSecretManager', 'GoogleCloudHsmGeneratedSecretManager'. If not set, no secret manager is used and the password is treated as a plain password.
usernamestrUsername for the JDBC source.
write_statementstrSQL query used to insert records into the JDBC sink.