MCPcopy Create free account

hub / github.com/aws-samples/amazon-managed-service-for-apache-flink-examples / types & classes

Types & classes133 in github.com/aws-samples/amazon-managed-service-for-apache-flink-examples

↓ 1 callersClassDeviceAggregation
python/DatastreamKafkaConnector/datastream-kafka-connector-example.py:49
Class
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.ts:32
Class
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.d.ts:3
Class
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.js:6
ClassAddress
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:189
ClassAggregateVehicleEvent
Record representing aggregate events from connected vehicle. <p> Similarly to {@link VehicleEvent}, it demonstrates the use of custom TypeInformation.
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:19
ClassAggregateVehicleEventSerializationTest
Tests the serialization of {@link AggregateVehicleEvent}
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/AggregateVehicleEventSerializationTest.java:16
ClassAggregateVehicleEventsWindowFunction
Implementation of ProcessWindowFunction that aggregates the sensor data and warning lights for a vehicle. <p> The actual implementation is not relevan
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/aggregation/AggregateVehicleEventsWindowFunction.java:21
ClassAggregatedStockPrice
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:5
ClassAverageSensorDataTypeInfoFactory
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:31
ClassAvroGenericRecordToRowDataMapper
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/AvroGenericRecordToRowDataMapper.java:14
ClassAvroGenericStockTradeGeneratorFunction
Function used by DataGen source to generate random records as AVRO GenericRecord. <p> The generator assumes that the AVRO schema provided contains the
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunction.java:17
ClassAvroGenericStockTradeGeneratorFunction
Function used by DataGen source to generate random records as AVRO GenericRecord. <p> The generator assumes that the AVRO schema provided contains the
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunction.java:17
ClassAvroGenericStockTradeGeneratorFunctionTest
java/Iceberg/S3TableSink/src/test/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunctionTest.java:10
ClassAvroSchemaUtils
java/Iceberg/IcebergDataStreamSource/src/main/java/com/amazonaws/services/msf/avro/AvroSchemaUtils.java:10
ClassAvroSchemaUtils
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/avro/AvroSchemaUtils.java:10
ClassAvroSchemaUtils
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/avro/AvroSchemaUtils.java:10
ClassAvroSpecificRecordBulkFormat
Simple BulkFormat to read AVRO files into SpecificRecord @param <O>
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/AvroSpecificRecordBulkFormat.java:16
ClassBasicStreamingJob
A basic Flink Java application to run on Amazon Managed Service for Apache Flink, with Kinesis Data Streams as source and sink.
java/GettingStarted/src/main/java/com/amazonaws/services/msf/BasicStreamingJob.java:23
ClassBasicTableJob
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/BasicTableJob.java:24
ClassChangeEvent
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:5
ClassChangeEventDeserializationSchema
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEventDeserializationSchema.java:12
ClassConfigurationHelper
java/JdbcSink/src/main/java/com/amazonaws/services/msf/ConfigurationHelper.java:9
ClassCustomTypeInfoJob
Basic streaming job demonstrating how to define a custom TypeInfo for your records
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/CustomTypeInfoJob.java:32
ClassDataGeneratorJob
A Flink application that generates random stock data using DataGeneratorSource and sends it to Kinesis Data Streams and/or Kafka as JSON based on conf
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:36
ClassDataGeneratorJobTest
java/FlinkDataGenerator/src/test/java/com/amazonaws/services/msf/DataGeneratorJobTest.java:17
ClassDdbTableItem
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:8
ClassFetchSecretsJob
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:31
ClassFirehoseStreamingJob
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/FirehoseStreamingJob.java:20
ClassFlinkCDCSqlServer2JdbcJob
java/FlinkCDC/FlinkCDCSQLServerSource/src/main/java/com/amazonaws/services/msf/FlinkCDCSqlServer2JdbcJob.java:17
ClassFlinkSerializationTestUtils
Utilities for testing Flink serialization.
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/FlinkSerializationTestUtils.java:19
ClassHadoopUtils
python/IcebergSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:9
ClassHadoopUtils
This class is a copy of org.apache.flink.runtime.util.HadoopUtils with the getHadoopConfiguration() method replaced to return an org.apache.hadoop.con
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:27
ClassHadoopUtils
This class is a copy of org.apache.flink.runtime.util.HadoopUtils with the getHadoopConfiguration() method replaced to return an org.apache.hadoop.con
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:27
ClassIcebergConfig
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:102
ClassIcebergSQLSinkJob
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:25
ClassIcebergSinkBuilder
Wraps the code to initialize an Iceberg sink that uses S3 Tables Internal catalog
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:25
ClassIcebergSinkBuilder
Wraps the code to initialize an Iceberg sink that uses Glue Data Catalog as catalog
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:34
ClassIncomingEvent
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:5
ClassIncomingEvent
java/SideOutputs/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:3
ClassIncomingEventDataGeneratorFunction
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEventDataGeneratorFunction.java:7
ClassIncomingEventDataGeneratorFunction
java/SideOutputs/src/main/java/com/amazonaws/services/msf/IncomingEventDataGeneratorFunction.java:6
ClassJdbcSinkJob
A Flink application that generates random stock price data using DataGeneratorSource and writes it to a PostgreSQL database using the JDBC connector.
java/JdbcSink/src/main/java/com/amazonaws/services/msf/JdbcSinkJob.java:33
ClassJsonConverter
Simple converter for Avro records to JSON. Not designed for production.
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:15
ClassJsonConverter
Simple converter for Avro records to JSON. Not designed for production.
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:15
ClassJsonSerializationTest
java/FlinkDataGenerator/src/test/java/com/amazonaws/services/msf/JsonSerializationTest.java:8
ClassKafkaStreamingJob
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:26
ClassKdaAutoscalingStack
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.ts:32
ClassKdaAutoscalingStack
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.d.ts:3
ClassKdaAutoscalingStack
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.js:6
ClassKeyHashKafkaPartitioner
Example of FlinkKafkaPartitioner which partitions by key. <p> The KafkaSink in DataStream API does not use the Kafka default partitioner by default.
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/kafka/KeyHashKafkaPartitioner.java:14
ClassKinesisDeaggregatingDeserializationSchemaWrapper
Wrapper adding de-aggregation to any DeserializationSchema @param <T> type returned by the DeserializationSchema
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:22
ClassKplAggregatingProducer
Simple KPL producer publishing random StockRecords as JSON to a Kinesis stream
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:27
ClassMetricEmitterNoOpMap
No-op Map function exposing 3 custom metrics: generatedRecordCount, generatedRecordRatePerParallelism, and taskParallelism. Note each subtask emits it
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/MetricEmitterNoOpMap.java:14
ClassMetricEmittingMapperFunction
NoOp mapper function which acts as a pass through for the records. It demonstrates reporting of custom metrics that get published to Cloud Watch by Ma
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/MetricEmittingMapperFunction.java:11
EnumMetricStats
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.ts:None
EnumMetricType
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.ts:None
ClassMinAggregate
Implementation of min aggregate for demonstration purposes. This specific outcome could be achieved more simply by relying on pre-defined aggregators.
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:123
ClassMoreKryoSerializationExamplesTest
This test collects different examples of record object that will or will not fall-back to Kryo for serialization. <p> Serialization tests are included
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:35
ClassNestedPojo
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:168
ClassPojoWithCollection
This POJO contains a Collection field. It falls back to Kryo unless you add a @TypeInfo
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:41
ClassPojoWithInstant
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:88
ClassPojoWithNonDirectlySerializableTimeFields
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:123
ClassProcessedEvent
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessedEvent.java:3
ClassProcessingFunction
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessingFunction.java:19
EnumProcessingOutcome
java/SideOutputs/src/main/java/com/amazonaws/services/msf/ProcessingOutcome.java:3
ClassRecordCountJob
A sample Managed Service For Apache Flink application with Kinesis data streams as source and sink with simple filter function.
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/RecordCountJob.java:25
ClassRecordDeaggregator
De-aggregate software.amazon.awssdk.services.kinesis.model.Record into a collection of software.amazon.kinesis.retrieval.KinesisClientRecord. <p> The
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:72
ClassRetriesFlinkJob
java/AsyncIO/src/main/java/com/amazonaws/services/msf/RetriesFlinkJob.java:30
ClassS3ParquetToKinesisJob
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/S3ParquetToKinesisJob.java:27
ClassSQSStreamingJob
java/SQSSink/src/main/java/com/amazonaws/services/msf/SQSStreamingJob.java:20
ClassSensorDataTypeInfoFactory
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:40
ClassSideOutputsFlinkJob
A sample Managed Service For Apache Flink application with Data Gen as a source and Kinesis data streams as the sinks demonstrating how to use Side Ou
java/SideOutputs/src/main/java/com/amazonaws/services/msf/SideOutputsFlinkJob.java:33
ClassSpeedLimitFilter
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/SpeedLimitFilter.java:3
ClassSpeedRecord
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/SpeedRecord.java:3
ClassSpeedRecordGeneratorFunction
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/SpeedRecordGeneratorFunction.java:7
ClassStock
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:5
ClassStockPrice
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:5
ClassStockPrice
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:3
ClassStockPrice
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:5
ClassStockPrice
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:5
ClassStockPrice
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:6
ClassStockPrice
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:3
ClassStockPrice
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:7
ClassStockPrice
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:5
ClassStockPrice
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:5
ClassStockPrice
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:5
ClassStockPrice
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:7
ClassStockPrice
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:3
ClassStockPriceGenerator
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPriceGenerator.java:8
ClassStockPriceGeneratorFunction
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
Generator function that creates random Stock objects. Implements GeneratorFunction to work with DataGeneratorSource. Modify this class to generate a
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPriceGeneratorFunction.java:15
ClassStockPriceGeneratorFunction
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/datagen/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
Function used by DataGen source to generate random records as com.amazonaws.services.msf.pojo.StockPrice POJOs. The generator mimics the behavior of
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/source/StockPriceGeneratorFunction.java:14
ClassStockPriceGeneratorFunction
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/datagen/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:9
ClassStockPriceGeneratorFunction
Generator function that creates realistic fake StockPrice objects using JavaFaker. Implements GeneratorFunction to work with DataGeneratorSource.
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPriceGeneratorFunction.java:17
next →1–100 of 133, ranked by callers