#
# 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.
#
# NOTE: This file contains autogenerated external transform(s)
# and should not be edited by hand.
# Refer to gen_xlang_wrappers.py for more info.
"""Cross-language transforms in this module can be imported from the
:py:mod:`apache_beam.io` package."""
# pylint:disable=line-too-long
from apache_beam.transforms.external import BeamJarExpansionService
from apache_beam.transforms.external_transform_provider import ExternalTransform
[docs]
class DatadogWrite(ExternalTransform):
identifier = "beam:schematransform:org.apache.beam:datadog_write:v1"
def __init__(
self,
api_key,
url,
batch_count=None,
error_handling=None,
max_buffer_size=None,
min_batch_count=None,
parallelism=None,
expansion_service=None):
"""
:param api_key: (str)
The Datadog API key.
:param url: (str)
The Datadog API URL.
:param batch_count: (int32)
The number of events to batch together for each write.
:param error_handling: (Row(output=<class 'str'>))
Specifies how to handle errors.
:param max_buffer_size: (int64)
The maximum buffer size in bytes.
:param min_batch_count: (int32)
The minimum number of events to batch together for each write.
:param parallelism: (int32)
The degree of parallelism for writing.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
api_key=api_key,
url=url,
batch_count=batch_count,
error_handling=error_handling,
max_buffer_size=max_buffer_size,
min_batch_count=min_batch_count,
parallelism=parallelism,
expansion_service=expansion_service)
[docs]
class GenerateSequence(ExternalTransform):
"""
Outputs a PCollection of Beam Rows, each containing a single INT64 number
called "value". The count is produced from the given "start" value and either
up to the given "end" or until 2^63 - 1.
To produce an unbounded PCollection, simply do not specify an "end" value.
Unbounded sequences can specify a "rate" for output elements.
In all cases, the sequence of numbers is generated in parallel, so there is no
inherent ordering between the generated values
"""
identifier = "beam:schematransform:org.apache.beam:generate_sequence:v1"
def __init__(self, start, end=None, rate=None, expansion_service=None):
"""
:param start: (int64)
The minimum number to generate (inclusive).
:param end: (int64)
The maximum number to generate (exclusive). Will be an unbounded
sequence if left unspecified.
:param rate: (Row(elements=<class 'int64'>, seconds=typing.Optional[int64]))
Specifies the rate to generate a given number of elements per a given
number of seconds. Applicable only to unbounded sequences.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
start=start, end=end, rate=rate, expansion_service=expansion_service)
[docs]
class MongodbRead(ExternalTransform):
identifier = "beam:schematransform:org.apache.beam:mongodb_read:v1"
def __init__(
self,
collection,
database,
schema,
uri,
error_handling=None,
filter=None,
expansion_service=None):
"""
:param collection: (str)
The MongoDB collection to read from.
:param database: (str)
The MongoDB database to read from.
:param schema: (str)
The schema in which the data is encoded, defined with JSON-schema syntax
(https://json-schema.org/).
:param uri: (str)
The connection URI for the MongoDB server.
:param error_handling: (Row(output=<class 'str'>))
This option specifies whether and where to output rows that failed to be
read.
:param filter: (str)
An optional BSON filter to apply to the read. This should be a valid
JSON string.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
collection=collection,
database=database,
schema=schema,
uri=uri,
error_handling=error_handling,
filter=filter,
expansion_service=expansion_service)
[docs]
class MongodbWrite(ExternalTransform):
identifier = "beam:schematransform:org.apache.beam:mongodb_write:v1"
def __init__(
self,
collection,
database,
uri,
batch_size=None,
error_handling=None,
expansion_service=None):
"""
:param collection: (str)
The MongoDB collection to write to.
:param database: (str)
The MongoDB database to write to.
:param uri: (str)
The connection URI for the MongoDB server.
:param batch_size: (int64)
The number of documents to include in each batch write.
:param error_handling: (Row(output=<class 'str'>))
This option specifies whether and where to output unwritable rows. Note:
Error handling is currently limited to data conversion failures before
sending to the MongoDB driver, as the underlying MongoDbIO does not yet
support dead-letter queues for write failures.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
collection=collection,
database=database,
uri=uri,
batch_size=batch_size,
error_handling=error_handling,
expansion_service=expansion_service)
[docs]
class KinesisRead(ExternalTransform):
"""
Reads records from an Amazon Kinesis stream and outputs Beam Rows with `data`
(bytes) plus stream metadata fields.
"""
identifier = "beam:schematransform:org.apache.beam:kinesis_read:v1"
def __init__(
self,
aws_access_key,
aws_secret_key,
region,
stream_name,
initial_position_in_stream=None,
initial_timestamp_in_stream=None,
max_capacity_per_shard=None,
max_num_records=None,
max_read_time=None,
rate_limit=None,
request_records_limit=None,
service_endpoint=None,
up_to_date_threshold=None,
verify_certificate=None,
watermark_idle_duration_threshold=None,
watermark_policy=None,
expansion_service=None):
"""
:param aws_access_key: (str)
AWS access key id.
:param aws_secret_key: (str)
AWS secret access key.
:param region: (str)
AWS region, for example us-east-1.
:param stream_name: (str)
Kinesis stream name to read from.
:param initial_position_in_stream: (str)
Where to start reading in the stream: LATEST, TRIM_HORIZON, or
AT_TIMESTAMP.
:param initial_timestamp_in_stream: (int64)
Epoch millis timestamp used when initial_position_in_stream is
AT_TIMESTAMP.
:param max_capacity_per_shard: (int64)
Maximum number of records to hold in memory per shard.
:param max_num_records: (int64)
Maximum number of records to read. When set, the resulting PCollection
is bounded.
:param max_read_time: (int64)
Maximum read time in milliseconds. When set, the resulting PCollection
is bounded.
:param rate_limit: (int64)
Fixed delay between GetRecords calls in milliseconds.
:param request_records_limit: (int64)
Max records returned by a single GetRecords call (1-10000).
:param service_endpoint: (str)
Optional custom Kinesis service endpoint URI.
:param up_to_date_threshold: (int64)
Threshold duration in milliseconds after which a shard is considered up
to date.
:param verify_certificate: (boolean)
Whether to verify TLS certificates. Defaults to true. Never disable in
production.
:param watermark_idle_duration_threshold: (int64)
Idle duration threshold in milliseconds for ARRIVAL_TIME watermark
policy.
:param watermark_policy: (str)
Watermark policy: ARRIVAL_TIME or PROCESSING_TIME.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
aws_access_key=aws_access_key,
aws_secret_key=aws_secret_key,
region=region,
stream_name=stream_name,
initial_position_in_stream=initial_position_in_stream,
initial_timestamp_in_stream=initial_timestamp_in_stream,
max_capacity_per_shard=max_capacity_per_shard,
max_num_records=max_num_records,
max_read_time=max_read_time,
rate_limit=rate_limit,
request_records_limit=request_records_limit,
service_endpoint=service_endpoint,
up_to_date_threshold=up_to_date_threshold,
verify_certificate=verify_certificate,
watermark_idle_duration_threshold=watermark_idle_duration_threshold,
watermark_policy=watermark_policy,
expansion_service=expansion_service)
[docs]
class KinesisWrite(ExternalTransform):
"""
Writes Beam Rows to an Amazon Kinesis stream. Each input row must include a
`data` (bytes) field, or a `payload` string/bytes field. Optional per-row
`partition_key` overrides the configured default partition key.
"""
identifier = "beam:schematransform:org.apache.beam:kinesis_write:v1"
def __init__(
self,
aws_access_key,
aws_secret_key,
partition_key,
region,
stream_name,
aggregation_enabled=None,
aggregation_max_buffered_time=None,
aggregation_max_bytes=None,
aggregation_shard_refresh_interval=None,
service_endpoint=None,
verify_certificate=None,
expansion_service=None):
"""
:param aws_access_key: (str)
AWS access key id.
:param aws_secret_key: (str)
AWS secret access key.
:param partition_key: (str)
Default partition key used when the input row does not contain a
partition_key field.
:param region: (str)
AWS region, for example us-east-1.
:param stream_name: (str)
Kinesis stream name to write to.
:param aggregation_enabled: (boolean)
Enable KPL-compatible record aggregation.
:param aggregation_max_buffered_time: (int64)
Max time in milliseconds to buffer records when aggregation is enabled.
:param aggregation_max_bytes: (int64)
Max aggregated record size in bytes when aggregation is enabled.
:param aggregation_shard_refresh_interval: (int64)
Shard map refresh interval in minutes when aggregation is enabled.
:param service_endpoint: (str)
Optional custom Kinesis service endpoint URI.
:param verify_certificate: (boolean)
Whether to verify TLS certificates. Defaults to true. Never disable in
production.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
aws_access_key=aws_access_key,
aws_secret_key=aws_secret_key,
partition_key=partition_key,
region=region,
stream_name=stream_name,
aggregation_enabled=aggregation_enabled,
aggregation_max_buffered_time=aggregation_max_buffered_time,
aggregation_max_bytes=aggregation_max_bytes,
aggregation_shard_refresh_interval=aggregation_shard_refresh_interval,
service_endpoint=service_endpoint,
verify_certificate=verify_certificate,
expansion_service=expansion_service)
[docs]
class TfrecordRead(ExternalTransform):
identifier = "beam:schematransform:org.apache.beam:tfrecord_read:v1"
def __init__(
self,
compression,
file_pattern,
validate,
error_handling=None,
expansion_service=None):
"""
:param compression: (str)
Decompression type to use when reading input files.
:param file_pattern: (str)
Filename or file pattern used to find input files.
:param validate: (boolean)
Validate file pattern.
:param error_handling: (Row(output=<class 'str'>))
This option specifies whether and where to output unwritable rows.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
compression=compression,
file_pattern=file_pattern,
validate=validate,
error_handling=error_handling,
expansion_service=expansion_service)
[docs]
class TfrecordWrite(ExternalTransform):
identifier = "beam:schematransform:org.apache.beam:tfrecord_write:v1"
def __init__(
self,
compression,
num_shards,
output_prefix,
error_handling=None,
filename_suffix=None,
max_num_writers_per_bundle=None,
no_spilling=None,
shard_template=None,
expansion_service=None):
"""
:param compression: (str)
Option to indicate the output sink's compression type. Default is NONE.
:param num_shards: (int32)
The number of shards to use, or 0 for automatic.
:param output_prefix: (str)
The directory to which files will be written.
:param error_handling: (Row(output=<class 'str'>))
This option specifies whether and where to output unwritable rows.
:param filename_suffix: (str)
The suffix of each file written, combined with prefix and shardTemplate.
:param max_num_writers_per_bundle: (int32)
Maximum number of writers created in a bundle before spilling to
shuffle.
:param no_spilling: (boolean)
Whether to skip the spilling of data caused by having
maxNumWritersPerBundle.
:param shard_template: (str)
The shard template of each file written, combined with prefix and
suffix.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:expansion-service:shadowJar")
super().__init__(
compression=compression,
num_shards=num_shards,
output_prefix=output_prefix,
error_handling=error_handling,
filename_suffix=filename_suffix,
max_num_writers_per_bundle=max_num_writers_per_bundle,
no_spilling=no_spilling,
shard_template=shard_template,
expansion_service=expansion_service)
[docs]
class ReadFromJms(ExternalTransform):
"""
Reads messages from a JMS broker and outputs each message as a single
`payload` (string) field.
By default the read is unbounded (streaming): it keeps consuming messages from
the specified queue or topic until the pipeline is stopped. Setting
`maxNumRecords` and/or `maxReadTimeSeconds` bounds the read, producing a
bounded (batch) PCollection.
"""
identifier = "beam:schematransform:org.apache.beam:jms_read:v1"
def __init__(
self,
connection_configuration,
acknowledge_mode=None,
close_timeout_seconds=None,
individual_acknowledge_mode_code=None,
max_num_records=None,
max_read_time_seconds=None,
queue=None,
topic=None,
expansion_service=None):
"""
:param connection_configuration: (Row(connection_factory_class_name=typing.Optional[str], password=typing.Optional[str], server_uri=<class 'str'>, username=typing.Optional[str]))
Configuration options to set up the JMS connection.
Note: if connection factory class name is set, spin up a persistent
expansion service with the provider client JAR on the classpath:
java -cp <expansion-service-jar>:<provider-client-jar>
org.apache.beam.sdk.expansion.service.ExpansionService <port>
and pass expansion_service='localhost:<port>' to the transform.
Currently, only ActiveMQ (org.apache.activemq.ActiveMQConnectionFactory)
is embedded into the expansion service
:param acknowledge_mode: (str)
The JMS acknowledge mode: CLIENT_ACKNOWLEDGE, CLIENT_ACKNOWLEDGE_UNSAFE,
or INDIVIDUAL_ACKNOWLEDGE.
:param close_timeout_seconds: (int64)
Close timeout for the JMS connection in seconds.
:param individual_acknowledge_mode_code: (int32)
The proprietary integer code for individual message acknowledgment when
using INDIVIDUAL_ACKNOWLEDGE.
:param max_num_records: (int64)
The max number of records to receive. Setting this will result in a
bounded PCollection.
:param max_read_time_seconds: (int64)
The maximum time for this source to read messages. Setting this will
result in a bounded PCollection.
:param queue: (str)
The JMS queue to read from. Exclusively one of queue or topic must be
specified.
:param topic: (str)
The JMS topic to read from. Exclusively one of queue or topic must be
specified.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:messaging-expansion-service:shadowJar")
super().__init__(
connection_configuration=connection_configuration,
acknowledge_mode=acknowledge_mode,
close_timeout_seconds=close_timeout_seconds,
individual_acknowledge_mode_code=individual_acknowledge_mode_code,
max_num_records=max_num_records,
max_read_time_seconds=max_read_time_seconds,
queue=queue,
topic=topic,
expansion_service=expansion_service)
[docs]
class WriteToJms(ExternalTransform):
"""
Publishes messages to a JMS broker. Expects an input PCollection of rows with
a `payload` (string) or `bytes` field, each of which is published as one JMS
TextMessage.
Works with both bounded (batch) and unbounded (streaming) input PCollections.
"""
identifier = "beam:schematransform:org.apache.beam:jms_write:v1"
def __init__(
self,
connection_configuration,
queue=None,
topic=None,
expansion_service=None):
"""
:param connection_configuration: (Row(connection_factory_class_name=typing.Optional[str], password=typing.Optional[str], server_uri=<class 'str'>, username=typing.Optional[str]))
Configuration options to set up the JMS connection.
Note: if connection factory class name is set, spin up a persistent
expansion service with the provider client JAR on the classpath:
java -cp <expansion-service-jar>:<provider-client-jar>
org.apache.beam.sdk.expansion.service.ExpansionService <port>
and pass expansion_service='localhost:<port>' to the transform.
Currently, only ActiveMQ (org.apache.activemq.ActiveMQConnectionFactory)
is embedded into the expansion service
:param queue: (str)
The JMS queue to write to. Exclusively one of queue or topic must be
specified.
:param topic: (str)
The JMS topic to write to. Exclusively one of queue or topic must be
specified.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:messaging-expansion-service:shadowJar")
super().__init__(
connection_configuration=connection_configuration,
queue=queue,
topic=topic,
expansion_service=expansion_service)
[docs]
class ReadFromMqtt(ExternalTransform):
"""
Reads messages from an MQTT broker and outputs each payload as a single
`bytes` field.
By default the read is unbounded (streaming): it keeps consuming messages from
the subscribed topic until the pipeline is stopped. Setting `maxNumRecords`
and/or `maxReadTimeSeconds` bounds the read, producing a bounded (batch)
PCollection.
Note: streaming reads require a runner that supports portable streaming (e.g.
Prism, Flink, or Dataflow). The legacy local Python DirectRunner does not
execute portable streaming cross-language reads.
"""
identifier = "beam:schematransform:org.apache.beam:mqtt_read:v1"
def __init__(
self,
connection_configuration,
max_num_records=None,
max_read_time_seconds=None,
expansion_service=None):
"""
:param connection_configuration: (Row(client_id=typing.Optional[str], password=typing.Optional[str], server_uri=<class 'str'>, topic=typing.Optional[str], username=typing.Optional[str]))
Configuration options to set up the MQTT connection.
:param max_num_records: (int64)
The max number of records to receive. Setting this will result in a
bounded PCollection.
:param max_read_time_seconds: (int64)
The maximum time for this source to read messages. Setting this will
result in a bounded PCollection.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:messaging-expansion-service:shadowJar")
super().__init__(
connection_configuration=connection_configuration,
max_num_records=max_num_records,
max_read_time_seconds=max_read_time_seconds,
expansion_service=expansion_service)
[docs]
class WriteToMqtt(ExternalTransform):
"""
Publishes messages to an MQTT broker. Expects an input PCollection of rows
with a single `bytes` field, each of which is published as one MQTT message.
Works with both bounded (batch) and unbounded (streaming) input PCollections.
"""
identifier = "beam:schematransform:org.apache.beam:mqtt_write:v1"
def __init__(
self, connection_configuration, retained=None, expansion_service=None):
"""
:param connection_configuration: (Row(client_id=typing.Optional[str], password=typing.Optional[str], server_uri=<class 'str'>, topic=typing.Optional[str], username=typing.Optional[str]))
Configuration options to set up the MQTT connection.
:param retained: (boolean)
Whether or not the publish message should be retained by the messaging
engine. When a subscriber connects, it gets the latest retained message.
Defaults to `False`, which will clear the retained message from the
server.
"""
self.default_expansion_service = BeamJarExpansionService(
"sdks:java:io:messaging-expansion-service:shadowJar")
super().__init__(
connection_configuration=connection_configuration,
retained=retained,
expansion_service=expansion_service)