Skip to content

Creates a Channel that streams records from an Amazon MSK Express cluster topic to Amazon S3 or Apache Iceberg

Description

Creates a Channel that streams records from an Amazon MSK Express cluster topic to Amazon S3 or Apache Iceberg.

Usage

kafka_create_channel(ChannelName, ClusterArn, EncryptionConfiguration,
  IcebergDestinationConfiguration, S3DestinationConfiguration, Tags,
  TopicConfigurationList, LoggingInfo)

Arguments

  • ChannelName

    [required] The name of the channel. Must be unique within the cluster.

  • ClusterArn

    [required] The Amazon Resource Name (ARN) that uniquely identifies the cluster.

  • EncryptionConfiguration

    The encryption configuration applied to the channel.

  • IcebergDestinationConfiguration

    The Apache Iceberg destination for the channel. Mutually exclusive with s3DestinationConfiguration.

  • S3DestinationConfiguration

    The Amazon S3 destination for the channel. Mutually exclusive with icebergDestinationConfiguration.

  • Tags

    The tags attached to the channel.

  • TopicConfigurationList

    [required] The list of topic configurations for the channel. Currently exactly one topic must be specified.

  • LoggingInfo

    The destinations to which the channel publishes operational logs.

Value

A list with the following syntax:

list(
  ChannelArn = "string",
  ClusterOperationArn = "string"
)

Request syntax

svc$create_channel(
  ChannelName = "string",
  ClusterArn = "string",
  EncryptionConfiguration = list(
    KmsKeyArn = "string"
  ),
  IcebergDestinationConfiguration = list(
    AppendOnly = TRUE|FALSE,
    Catalog = list(
      CatalogArn = "string",
      WarehouseLocation = "string"
    ),
    DataFreshnessInSeconds = 123,
    DeadLetterQueueS3 = list(
      BucketArn = "string",
      ErrorOutputPrefix = "string",
      ExpectedBucketOwner = "string"
    ),
    DestinationTableList = list(
      list(
        DestinationDatabaseName = "string",
        DestinationTableName = "string",
        PartitionSpec = list(
          PartitionStrategy = "TIME_HOUR",
          SourceList = list(
            list(
              SourceName = "string"
            )
          )
        )
      )
    ),
    SchemaEvolution = list(
      EnableSchemaEvolution = TRUE|FALSE
    ),
    ServiceExecutionRoleArn = "string",
    TableCreation = list(
      EnableTableCreation = TRUE|FALSE
    ),
    CompressionType = "ZSTD"|"SNAPPY"
  ),
  S3DestinationConfiguration = list(
    DataFreshnessInSeconds = 123,
    DeadLetterQueueS3 = list(
      BucketArn = "string",
      ErrorOutputPrefix = "string",
      ExpectedBucketOwner = "string"
    ),
    ServiceExecutionRoleArn = "string",
    Storage = list(
      BucketArn = "string",
      CompressionType = "NONE"|"GZIP"|"ZSTD",
      OutputPrefix = "string",
      OutputKeyTemplate = "string",
      StorageClass = "STANDARD"|"INTELLIGENT_TIERING"|"GLACIER_IR",
      ExpectedBucketOwner = "string"
    )
  ),
  Tags = list(
    "string"
  ),
  TopicConfigurationList = list(
    list(
      RecordConverter = list(
        ValueConverter = "BYTE_ARRAY"|"JSON"|"JSON_SCHEMA_GSR"|"STRING"
      ),
      RecordSchema = list(
        GsrArn = "string"
      ),
      TopicArn = "string"
    )
  ),
  LoggingInfo = list(
    CloudWatchLogs = list(
      Enabled = TRUE|FALSE,
      LogGroup = "string"
    ),
    Firehose = list(
      DeliveryStream = "string",
      Enabled = TRUE|FALSE
    ),
    S3 = list(
      Bucket = "string",
      Enabled = TRUE|FALSE,
      Prefix = "string"
    )
  )
)