blob: 173236dc5e942efba34d6edec9440fbe424ba033 [file]
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
# This file enumerates the various IOs that are available by default as
# top-level transforms in Beam's YAML.
#
# Note that there may be redundant implementations. In these cases the specs
# should be kept in sync.
# TODO(yaml): See if this can be enforced programmatically.
# BigQuery
- type: renaming
transforms:
'ReadFromBigQuery': 'ReadFromBigQuery'
'WriteToBigQuery': 'WriteToBigQuery'
config:
mappings:
'ReadFromBigQuery':
query: 'query'
table: 'table_spec'
fields: 'selected_fields'
row_restriction: 'row_restriction'
'WriteToBigQuery':
table: 'table'
create_disposition: 'create_disposition'
write_disposition: 'write_disposition'
error_handling: 'error_handling'
# TODO(https://github.com/apache/beam/issues/30058): Required until autosharding support is fixed
num_streams: 'num_streams'
underlying_provider:
type: beamJar
transforms:
'ReadFromBigQuery': 'beam:schematransform:org.apache.beam:bigquery_storage_read:v1'
'WriteToBigQuery': 'beam:schematransform:org.apache.beam:bigquery_write:v1'
config:
gradle_target: 'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
managed_replacement:
# Following transforms may be replaced with equivalent managed transforms,
# if the pipelines 'updateCompatibilityBeamVersion' match the provided
# version.
'ReadFromBigQuery': '2.69.0'
'WriteToBigQuery': '2.69.0'
# Kafka
- type: renaming
transforms:
'ReadFromKafka': 'ReadFromKafka'
'WriteToKafka': 'WriteToKafka'
config:
mappings:
'ReadFromKafka':
'schema': 'schema'
'consumer_config': 'consumer_config_updates'
'format': 'format'
'topic': 'topic'
'bootstrap_servers': 'bootstrap_servers'
'confluent_schema_registry_url': 'confluent_schema_registry_url'
'confluent_schema_registry_subject': 'confluent_schema_registry_subject'
'auto_offset_reset_config': 'auto_offset_reset_config'
'error_handling': 'error_handling'
'file_descriptor_path': 'file_descriptor_path'
'message_name': 'message_name'
'max_read_time_seconds': 'max_read_time_seconds'
'WriteToKafka':
'format': 'format'
'topic': 'topic'
'bootstrap_servers': 'bootstrap_servers'
'producer_config_updates': 'producer_config_updates'
'error_handling': 'error_handling'
'file_descriptor_path': 'file_descriptor_path'
'message_name': 'message_name'
'schema': 'schema'
underlying_provider:
type: beamJar
transforms:
'ReadFromKafka': 'beam:schematransform:org.apache.beam:kafka_read:v1'
'WriteToKafka': 'beam:schematransform:org.apache.beam:kafka_write:v1'
config:
gradle_target: 'sdks:java:io:expansion-service:shadowJar'
managed_replacement:
# Following transforms may be replaced with equivalent managed transforms,
# if the pipelines 'updateCompatibilityBeamVersion' match the provided
# version.
'ReadFromKafka': '2.65.0'
'WriteToKafka': '2.65.0'
# TODO(yaml): Tests are assuming python providers are before java ones, hence
# the order below. This should be fixed in the future.
# Python Providers
- type: python
transforms:
'ReadFromBigQuery': 'apache_beam.yaml.yaml_io.read_from_bigquery'
# Disable until https://github.com/apache/beam/issues/28162 is resolved.
# 'WriteToBigQuery': 'apache_beam.yaml.yaml_io.write_to_bigquery'
'ReadFromText': 'apache_beam.yaml.yaml_io.read_from_text'
'WriteToText': 'apache_beam.yaml.yaml_io.write_to_text'
'ReadFromPubSub': 'apache_beam.yaml.yaml_io.read_from_pubsub'
'WriteToPubSub': 'apache_beam.yaml.yaml_io.write_to_pubsub'
'ReadFromIceberg': 'apache_beam.yaml.yaml_io.read_from_iceberg'
'WriteToIceberg': 'apache_beam.yaml.yaml_io.write_to_iceberg'
'ReadFromTFRecord': 'apache_beam.yaml.yaml_io.read_from_tfrecord'
'WriteToTFRecord': 'apache_beam.yaml.yaml_io.write_to_tfrecord'
'ReadFromMongoDB': 'apache_beam.yaml.yaml_io.read_from_mongodb'
'WriteToMongoDB': 'apache_beam.yaml.yaml_io.write_to_mongodb'
'ReadFromDelta': 'apache_beam.yaml.yaml_io.read_from_delta'
'DicomSearch': 'apache_beam.yaml.yaml_io.dicom_search'
# General File Formats
# Declared as a renaming transform to avoid exposing all
# (implementation-specific) pandas arguments and aligning with possible Java
# implementation.
# Invoking these directly as a PyTransform is still an option for anyone wanting
# to use these power-features in a language-dependent manner.
- type: renaming
transforms:
'ReadFromCsv': 'ReadFromCsv'
'WriteToCsv': 'WriteToCsv'
'ReadFromJson': 'ReadFromJson'
'WriteToJson': 'WriteToJson'
'ReadFromParquet': 'ReadFromParquet'
'WriteToParquet': 'WriteToParquet'
'ReadFromAvro': 'ReadFromAvro'
'WriteToAvro': 'WriteToAvro'
config:
mappings:
'ReadFromCsv':
path: 'path'
delimiter: 'sep'
comment: 'comment'
filename_column: 'filename_column'
'WriteToCsv':
path: 'path'
delimiter: 'sep'
'ReadFromJson':
path: 'path'
'WriteToJson':
path: 'path'
'ReadFromParquet':
path: 'file_pattern'
'WriteToParquet':
path: 'file_path_prefix'
'ReadFromAvro':
path: 'file_pattern'
'WriteToAvro':
path: 'file_path_prefix'
defaults:
'ReadFromParquet':
as_rows: True
'ReadFromAvro':
as_rows: True
underlying_provider:
type: python
transforms:
'ReadFromCsv': 'apache_beam.io.ReadFromCsv'
'WriteToCsv': 'apache_beam.io.WriteToCsv'
'ReadFromJson': 'apache_beam.io.ReadFromJson'
'WriteToJson': 'apache_beam.io.WriteToJson'
'ReadFromParquet': 'apache_beam.io.ReadFromParquet'
'WriteToParquet': 'apache_beam.io.WriteToParquet'
'ReadFromAvro': 'apache_beam.io.ReadFromAvro'
'WriteToAvro': 'apache_beam.io.WriteToAvro'
# BeamJar Providers
- type: beamJar
transforms:
'WriteToCsv': 'beam:schematransform:org.apache.beam:csv_write:v1'
'WriteToJson': 'beam:schematransform:org.apache.beam:json_write:v1'
config:
gradle_target: 'sdks:java:extensions:schemaio-expansion-service:shadowJar'
# JMS
- type: renaming
transforms:
'ReadFromJms': 'ReadFromJms'
'WriteToJms': 'WriteToJms'
config:
mappings:
'ReadFromJms':
connection_configuration: 'connection_configuration'
queue: 'queue'
topic: 'topic'
max_num_records: 'max_num_records'
max_read_time_seconds: 'max_read_time_seconds'
close_timeout_seconds: 'close_timeout_seconds'
acknowledge_mode: 'acknowledge_mode'
individual_acknowledge_mode_code: 'individual_acknowledge_mode_code'
'WriteToJms':
connection_configuration: 'connection_configuration'
queue: 'queue'
topic: 'topic'
underlying_provider:
type: beamJar
transforms:
'ReadFromJms': 'beam:schematransform:org.apache.beam:jms_read:v1'
'WriteToJms': 'beam:schematransform:org.apache.beam:jms_write:v1'
config:
gradle_target: 'sdks:java:io:messaging-expansion-service:shadowJar'
- type: renaming
transforms:
'ReadFromIbmMQ': 'ReadFromIbmMQ'
'WriteToIbmMQ': 'WriteToIbmMQ'
config:
mappings:
'ReadFromIbmMQ':
connection_configuration: 'connection_configuration'
queue: 'queue'
topic: 'topic'
max_num_records: 'max_num_records'
max_read_time_seconds: 'max_read_time_seconds'
close_timeout_seconds: 'close_timeout_seconds'
acknowledge_mode: 'acknowledge_mode'
individual_acknowledge_mode_code: 'individual_acknowledge_mode_code'
'WriteToIbmMQ':
connection_configuration: 'connection_configuration'
queue: 'queue'
topic: 'topic'
underlying_provider:
type: beamJar
transforms:
'ReadFromIbmMQ': 'beam:schematransform:org.apache.beam:jms_read:v1'
'WriteToIbmMQ': 'beam:schematransform:org.apache.beam:jms_write:v1'
config:
gradle_target: 'sdks:java:io:messaging-expansion-service:shadowJar'
classpath:
- 'com.ibm.mq:com.ibm.mq.allclient:9.3.0.25'
- 'org.json:json:20251224'
# Debezium
- type: renaming
transforms:
'ReadFromDebezium': 'ReadFromDebezium'
config:
mappings:
'ReadFromDebezium':
connector: 'database'
username: 'username'
password: 'password'
host: 'host'
port: 'port'
table: 'table'
connection_properties: 'debezium_connection_properties'
max_number_of_records: 'max_number_of_records'
max_time_to_run: 'max_time_to_run'
underlying_provider:
type: beamJar
transforms:
'ReadFromDebezium': 'beam:schematransform:org.apache.beam:debezium_read:v1'
config:
gradle_target: 'sdks:java:io:debezium:expansion-service:shadowJar'
# Kinesis
- type: renaming
transforms:
'ReadFromKinesis': 'ReadFromKinesis'
'WriteToKinesis': 'WriteToKinesis'
config:
mappings:
'ReadFromKinesis':
stream_name: 'stream_name'
aws_access_key: 'aws_access_key'
aws_secret_key: 'aws_secret_key'
region: 'region'
service_endpoint: 'service_endpoint'
verify_certificate: 'verify_certificate'
max_num_records: 'max_num_records'
max_read_time: 'max_read_time'
initial_position_in_stream: 'initial_position_in_stream'
initial_timestamp_in_stream: 'initial_timestamp_in_stream'
request_records_limit: 'request_records_limit'
up_to_date_threshold: 'up_to_date_threshold'
max_capacity_per_shard: 'max_capacity_per_shard'
watermark_policy: 'watermark_policy'
watermark_idle_duration_threshold: 'watermark_idle_duration_threshold'
rate_limit: 'rate_limit'
'WriteToKinesis':
stream_name: 'stream_name'
aws_access_key: 'aws_access_key'
aws_secret_key: 'aws_secret_key'
region: 'region'
partition_key: 'partition_key'
service_endpoint: 'service_endpoint'
verify_certificate: 'verify_certificate'
aggregation_enabled: 'aggregation_enabled'
aggregation_max_bytes: 'aggregation_max_bytes'
aggregation_max_buffered_time: 'aggregation_max_buffered_time'
aggregation_shard_refresh_interval: 'aggregation_shard_refresh_interval'
underlying_provider:
type: beamJar
transforms:
'ReadFromKinesis': 'beam:schematransform:org.apache.beam:kinesis_read:v1'
'WriteToKinesis': 'beam:schematransform:org.apache.beam:kinesis_write:v1'
config:
gradle_target: 'sdks:java:io:amazon-web-services2:expansion-service:shadowJar'
# Snowflake
- type: renaming
transforms:
'ReadFromSnowflake': 'ReadFromSnowflake'
'WriteToSnowflake': 'WriteToSnowflake'
config:
mappings:
'ReadFromSnowflake':
server_name: 'server_name'
username: 'username'
password: 'password'
oauth_token: 'oauth_token'
private_key: 'private_key'
private_key_passphrase: 'private_key_passphrase'
database: 'database'
snowflake_schema: 'snowflake_schema'
warehouse: 'warehouse'
role: 'role'
table: 'table'
query: 'query'
staging_bucket_name: 'staging_bucket_name'
storage_integration_name: 'storage_integration_name'
schema: 'schema'
quotation_mark: 'quotation_mark'
'WriteToSnowflake':
server_name: 'server_name'
username: 'username'
password: 'password'
oauth_token: 'oauth_token'
private_key: 'private_key'
private_key_passphrase: 'private_key_passphrase'
database: 'database'
schema: 'schema'
warehouse: 'warehouse'
role: 'role'
table: 'table'
snow_pipe: 'snow_pipe'
staging_bucket_name: 'staging_bucket_name'
storage_integration_name: 'storage_integration_name'
create_disposition: 'create_disposition'
write_disposition: 'write_disposition'
quotation_mark: 'quotation_mark'
flush_row_limit: 'flush_row_limit'
flush_time_limit_millis: 'flush_time_limit_millis'
shards_number: 'shards_number'
debug_mode: 'debug_mode'
underlying_provider:
type: beamJar
transforms:
'ReadFromSnowflake':
'beam:schematransform:org.apache.beam:snowflake_read:v1'
'WriteToSnowflake':
'beam:schematransform:org.apache.beam:snowflake_write:v1'
config:
gradle_target: 'sdks:java:io:snowflake:expansion-service:shadowJar'
# Databases
- type: renaming
transforms:
'ReadFromJdbc': 'ReadFromJdbc'
'WriteToJdbc': 'WriteToJdbc'
'ReadFromMySql': 'ReadFromMySql'
'WriteToMySql': 'WriteToMySql'
'ReadFromPostgres': 'ReadFromPostgres'
'WriteToPostgres': 'WriteToPostgres'
'ReadFromOracle': 'ReadFromOracle'
'WriteToOracle': 'WriteToOracle'
'ReadFromSqlServer': 'ReadFromSqlServer'
'WriteToSqlServer': 'WriteToSqlServer'
config:
mappings:
'ReadFromJdbc':
url: 'jdbc_url'
connection_init_sql: 'connection_init_sql'
connection_properties: 'connection_properties'
disable_auto_commit: 'disable_auto_commit'
driver_class_name: 'driver_class_name'
driver_jars: 'driver_jars'
fetch_size: 'fetch_size'
output_parallelization: 'output_parallelization'
password: 'password'
query: 'read_query'
table: 'location'
partition_column : 'partition_column'
num_partitions: 'num_partitions'
secret_manager: 'secret_manager'
type: 'jdbc_type'
username: 'username'
'WriteToJdbc':
url: 'jdbc_url'
auto_sharding: 'autosharding'
connection_init_sql: 'connection_init_sql'
connection_properties: 'connection_properties'
driver_class_name: 'driver_class_name'
driver_jars: 'driver_jars'
password: 'password'
table: 'location'
batch_size: 'batch_size'
secret_manager: 'secret_manager'
type: 'jdbc_type'
username: 'username'
query: 'write_statement'
'ReadFromMySql': 'ReadFromJdbc'
'WriteToMySql': 'WriteToJdbc'
'ReadFromPostgres': 'ReadFromJdbc'
'WriteToPostgres': 'WriteToJdbc'
'ReadFromOracle': 'ReadFromJdbc'
'WriteToOracle': 'WriteToJdbc'
'ReadFromSqlServer': 'ReadFromJdbc'
'WriteToSqlServer': 'WriteToJdbc'
defaults:
'ReadFromMySql':
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'WriteToMySql':
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'ReadFromPostgres':
connection_init_sql: []
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'WriteToPostgres':
connection_init_sql: []
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'ReadFromOracle':
connection_init_sql: []
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'WriteToOracle':
connection_init_sql: []
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'ReadFromSqlServer':
connection_init_sql: []
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
'WriteToSqlServer':
connection_init_sql: []
driver_class_name: ''
driver_jars: ''
jdbc_type: ''
underlying_provider:
type: beamJar
transforms:
'ReadFromJdbc': 'beam:schematransform:org.apache.beam:jdbc_read:v1'
'ReadFromMySql': 'beam:schematransform:org.apache.beam:mysql_read:v1'
'ReadFromPostgres': 'beam:schematransform:org.apache.beam:postgres_read:v1'
'ReadFromOracle': 'beam:schematransform:org.apache.beam:oracle_read:v1'
'ReadFromSqlServer': 'beam:schematransform:org.apache.beam:sql_server_read:v1'
'WriteToJdbc': 'beam:schematransform:org.apache.beam:jdbc_write:v1'
'WriteToMySql': 'beam:schematransform:org.apache.beam:mysql_write:v1'
'WriteToPostgres': 'beam:schematransform:org.apache.beam:postgres_write:v1'
'WriteToOracle': 'beam:schematransform:org.apache.beam:oracle_write:v1'
'WriteToSqlServer': 'beam:schematransform:org.apache.beam:sql_server_write:v1'
config:
gradle_target: 'sdks:java:extensions:schemaio-expansion-service:shadowJar'
managed_replacement:
# Following transforms may be replaced with equivalent managed transforms,
# if the pipelines 'updateCompatibilityBeamVersion' match the provided
# version.
'ReadFromPostgres': '2.73.0'
'WriteToPostgres': '2.73.0'
'ReadFromMySql': '2.73.0'
'WriteToMySql': '2.73.0'
'ReadFromSqlServer': '2.73.0'
'WriteToSqlServer': '2.73.0'
# Spanner
- type: renaming
transforms:
'ReadFromSpanner': 'ReadFromSpanner'
'WriteToSpanner': 'WriteToSpanner'
config:
mappings:
'ReadFromSpanner':
project: 'project_id'
instance: 'instance_id'
database: 'database_id'
table: 'table_id'
query: 'query'
columns: 'columns'
index: 'index'
batching: 'batching'
error_handling: 'error_handling'
'WriteToSpanner':
project: 'project_id'
instance: 'instance_id'
database: 'database_id'
table: 'table_id'
error_handling: 'error_handling'
underlying_provider:
type: beamJar
transforms:
'ReadFromSpanner': 'beam:schematransform:org.apache.beam:spanner_read:v1'
'WriteToSpanner': 'beam:schematransform:org.apache.beam:spanner_write:v1'
config:
gradle_target: 'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
# Firestore
- type: renaming
transforms:
'ReadFromFirestore': 'ReadFromFirestore'
'WriteToFirestore': 'WriteToFirestore'
config:
mappings:
'ReadFromFirestore':
project: 'project_id'
database: 'database_id'
collection: 'collection_id'
schema: 'schema'
error_handling: 'error_handling'
'WriteToFirestore':
project: 'project_id'
database: 'database_id'
collection: 'collection_id'
document_id_field: 'document_id_field'
error_handling: 'error_handling'
underlying_provider:
type: beamJar
transforms:
'ReadFromFirestore': 'beam:schematransform:org.apache.beam:firestore_read:v1'
'WriteToFirestore': 'beam:schematransform:org.apache.beam:firestore_write:v1'
config:
gradle_target: 'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
# TFRecord
- type: renaming
transforms:
'ReadFromTFRecord': 'ReadFromTFRecord'
'WriteToTFRecord': 'WriteToTFRecord'
config:
mappings:
'ReadFromTFRecord':
file_pattern: 'file_pattern'
compression_type: 'compression'
validate: 'validate'
'WriteToTFRecord':
file_path_prefix: 'output_prefix'
shard_name_template: 'shard_template'
file_name_suffix: 'filename_suffix'
num_shards: 'num_shards'
compression_type: 'compression'
underlying_provider:
type: beamJar
transforms:
'ReadFromTFRecord': 'beam:schematransform:org.apache.beam:tfrecord_read:v1'
'WriteToTFRecord': 'beam:schematransform:org.apache.beam:tfrecord_write:v1'
config:
gradle_target: 'sdks:java:io:expansion-service:shadowJar'
#BigTable
- type: renaming
transforms:
'ReadFromBigTable': 'ReadFromBigTable'
'WriteToBigTable': 'WriteToBigTable'
config:
mappings:
#Temp removing read from bigTable IO
'ReadFromBigTable':
project: 'project_id'
instance: 'instance_id'
table: 'table_id'
flatten: "flatten"
'WriteToBigTable':
project: 'project_id'
instance: 'instance_id'
table: 'table_id'
underlying_provider:
type: beamJar
transforms:
'ReadFromBigTable': 'beam:schematransform:org.apache.beam:bigtable_read:v1'
'WriteToBigTable': 'beam:schematransform:org.apache.beam:bigtable_write:v1'
config:
gradle_target: 'sdks:java:io:google-cloud-platform:expansion-service:shadowJar'
#IcebergCDC
- type: renaming
transforms:
'ReadFromIcebergCDC': 'ReadFromIcebergCDC'
config:
mappings:
'ReadFromIcebergCDC':
table: 'table'
catalog_name: 'catalog_name'
catalog_properties: 'catalog_properties'
config_properties: 'config_properties'
drop: 'drop'
filter: 'filter'
from_snapshot: 'from_snapshot'
from_timestamp: 'from_timestamp'
keep: 'keep'
poll_interval_seconds: 'poll_interval_seconds'
starting_strategy: 'starting_strategy'
streaming: 'streaming'
to_snapshot: 'to_snapshot'
to_timestamp: 'to_timestamp'
underlying_provider:
type: beamJar
transforms:
'ReadFromIcebergCDC': 'beam:schematransform:org.apache.beam:iceberg_cdc_read:v1'
config:
gradle_target: 'sdks:java:io:expansion-service:shadowJar'
#IcebergAddFiles
- type: renaming
transforms:
'IcebergAddFiles': 'IcebergAddFiles'
config:
mappings:
'IcebergAddFiles':
table: 'table'
catalog_properties: 'catalog_properties'
config_properties: 'config_properties'
triggering_frequency_seconds: 'triggering_frequency_seconds'
manifest_file_size: 'manifest_file_size'
location_prefix: 'location_prefix'
partition_fields: 'partition_fields'
table_properties: 'table_properties'
sort_fields: 'sort_fields'
error_handling: 'error_handling'
underlying_provider:
type: beamJar
transforms:
'IcebergAddFiles': 'beam:schematransform:iceberg_add_files:v1'
config:
gradle_target: 'sdks:java:io:expansion-service:shadowJar'
#Datadog
- type: renaming
transforms:
'WriteToDatadog': 'WriteToDatadog'
config:
mappings:
'WriteToDatadog':
'url': 'url'
'api_key': 'api_key'
'min_batch_count': 'min_batch_count'
'batch_count': 'batch_count'
'max_buffer_size': 'max_buffer_size'
'parallelism': 'parallelism'
'error_handling': 'error_handling'
underlying_provider:
type: beamJar
transforms:
'WriteToDatadog': 'beam:schematransform:org.apache.beam:datadog_write:v1'
config:
gradle_target: 'sdks:java:io:expansion-service:shadowJar'
#MongoDB
- type: renaming
transforms:
'ReadFromMongoDB': 'ReadFromMongoDB'
'WriteToMongoDB': 'WriteToMongoDB'
config:
mappings:
'ReadFromMongoDB':
uri: "uri"
database: "database"
collection: "collection"
schema: "schema"
filter: "filter"
error_handling: "error_handling"
'WriteToMongoDB':
uri: "uri"
database: "database"
collection: "collection"
batch_size: "batch_size"
error_handling: "error_handling"
underlying_provider:
type: beamJar
transforms:
'ReadFromMongoDB': 'beam:schematransform:org.apache.beam:mongodb_read:v1'
'WriteToMongoDB': 'beam:schematransform:org.apache.beam:mongodb_write:v1'
config:
gradle_target: 'sdks:java:io:expansion-service:shadowJar'