Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/aws-samples/amazon-managed-service-for-apache-flink-examples
/ functions
Functions
544 in github.com/aws-samples/amazon-managed-service-for-apache-flink-examples
⨍
Functions
544
◇
Types & classes
133
↳
Endpoints
3
Method
main
(String[] args)
java/Iceberg/IcebergDataStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:60
Method
main
(String[] args)
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:71
Method
main
(String[] args)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:67
Method
main
(String[] args)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:89
Method
main
(String[] args)
java/SQSSink/src/main/java/com/amazonaws/services/msf/SQSStreamingJob.java:62
Method
main
(String[] args)
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/BasicTableJob.java:46
Method
main
(String[] args)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:70
Method
main
(String[] args)
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:74
Method
main
(String[] args)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/S3ParquetToKinesisJob.java:78
Method
main
(String[] args)
java/FlinkCDC/FlinkCDCSQLServerSource/src/main/java/com/amazonaws/services/msf/FlinkCDCSqlServer2JdbcJob.java:44
Method
main
(String[] args)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/FirehoseStreamingJob.java:63
Method
main
(String[] args)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/RetriesFlinkJob.java:47
Method
main
(String[] args)
java/GettingStarted/src/main/java/com/amazonaws/services/msf/BasicStreamingJob.java:61
Method
main
(String[] args)
java/KafkaConfigProviders/Kafka-SASL_SSL-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:148
Method
main
(String[] args)
java/KafkaConfigProviders/Kafka-mTLS-Keystore-Sql-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:59
Method
main
(String[] args)
java/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:105
Method
main
(String[] args)
java/SideOutputs/src/main/java/com/amazonaws/services/msf/SideOutputsFlinkJob.java:59
Method
main
(String[] args)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/JdbcSinkJob.java:131
Method
main
(String[] args)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:126
Method
main
(String[] args)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:95
Method
main
(String[] args)
java/AvroGlueSchemaRegistryKinesis/src/main/java/com/amazonaws/services/msf/StreamingJob.java:145
Method
main
(String[] args)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/CustomTypeInfoJob.java:68
Method
map
(Long aLong)
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/datagen/StockPriceGeneratorFunction.java:12
Method
map
(Long aLong)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/SpeedRecordGeneratorFunction.java:9
Method
map
(Long sequence)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/source/StockPriceGeneratorFunction.java:18
Method
map
(GenericRecord genericRecord)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/AvroGenericRecordToRowDataMapper.java:21
Method
map
Generates a random trade
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunction.java:30
Method
map
(Long aLong)
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:12
Method
map
(Long aLong)
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:12
Method
map
(Long value)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPriceGeneratorFunction.java:23
Method
map
(TemperatureSample value)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/map/TemperatureSampleToPrometheusTimeSeriesMapper.java:11
Method
map
(Long value)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/datagen/VehicleEventGeneratorFunction.java:18
Method
nestedPojoSerializeAndDeserializeWithoutKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:218
Method
nestedPojosShouldNotFallbackToKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:211
Method
onFailure
(Throwable t)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:173
Method
onSuccess
(UserRecordResult result)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:167
Method
open
(self, runtime_context: RuntimeContext)
python/DatastreamKafkaConnector/datastream-kafka-connector-example.py:51
Method
open
(DeserializationSchema.InitializationContext context)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:32
Method
open
(Configuration config)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/MetricEmittingMapperFunction.java:22
Method
open
(OpenContext openContext)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:24
Method
open
Instantiate the connection to an async client here to use within asyncInvoke
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessingFunction.java:40
Method
partition
(T record, byte[] key, byte[] value, String targetTopic, int[] partitions)
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/kafka/KeyHashKafkaPartitioner.java:16
Method
pojoWithCollectionShouldNotFallbackToKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:75
Method
pojoWithInstantShouldNotFallbackToKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:101
Method
pojoWithInstantShouldSerializeAndDeserializeWithoutKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:107
Method
pojoWithZonedDateTimeShouldNotFallbackToKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:156
Method
process
(String key, Context context, Iterable<Double> average
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:156
Method
processElement
(Tuple2<IncomingEvent, ProcessingOutcome> value, Context ctx, Collector<IncomingEv
java/SideOutputs/src/main/java/com/amazonaws/services/msf/SideOutputsFlinkJob.java:81
Method
query
Returns the SQL of the PreparedStatement @return SQL
java/JdbcSink/src/main/java/com/amazonaws/services/msf/jdbc/StockPriceUpsertQueryStatement.java:74
Function
removePrefix
(str: string, prefix: string)
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.ts:12
Method
serializationShouldNotFallbackToKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/VehicleEventSerializationTest.java:19
Method
serializationShouldNotFallbackToKryo
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/AggregateVehicleEventSerializationTest.java:18
Method
setAverageSensorData
(Map<String, Long> averageSensorData)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:71
Method
setCountDistinctWarnings
(int countDistinctWarnings)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:79
Method
setDuration
(Duration duration)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:140
Method
setElements
(List<String> elements)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:53
Method
setEventTime
(String eventTime)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:19
Method
setEventTime
(String eventTime)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:18
Method
setEventTime
(Timestamp eventTime)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:24
Method
setEventTime
(Timestamp eventTime)
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:24
Method
setEventTime
(String eventTime)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:19
Method
setEventTime
(String eventTime)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:27
Method
setEventTime
(Timestamp eventTime)
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:24
Method
setEventTime
(Timestamp eventTime)
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:24
Method
setEventTime
(Timestamp eventTime)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:24
Method
setFieldOne
(String fieldOne)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:41
Method
setFieldTwo
(String fieldTwo)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:49
Method
setLocalDateTime
(LocalDateTime localDateTime)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:149
Method
setMapOfElements
(Map<String, Integer> mapOfElements)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:61
Method
setMinPrice
(Double minPrice)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:40
Method
setPartitionKey
(String partitionKey)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:23
Method
setPrice
(float price)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:35
Method
setPrice
(double price)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:34
Method
setPrice
(Double price)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:40
Method
setPrice
(Double price)
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:40
Method
setPrice
(float price)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:35
Method
setPrice
(Float price)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:39
Method
setPrice
(float price)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:43
Method
setPrice
(double price)
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:40
Method
setPrice
(Double price)
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:40
Method
setPrice
(Double price)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:40
Method
setPrice
(BigDecimal price)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:42
Method
setPrice
(double price)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:36
Method
setRoomId
(String roomId)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:23
Method
setSensorData
(Map<String, Long> sensorData)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:106
Method
setSetOfElements
(Set<String> setOfElements)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:69
Method
setSortKey
(String sortKey)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:33
Method
setSymbol
(String symbol)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:31
Method
setSymbol
(String symbol)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:26
Method
setSymbol
(String symbol)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:20
Method
setTicker
(String ticker)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:27
Method
setTicker
(String ticker)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:26
Method
setTicker
(String ticker)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:32
Method
setTicker
(String ticker)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:24
Method
setTicker
(String ticker)
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:32
Method
setTicker
(String ticker)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:27
Method
setTicker
(String ticker)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:35
Method
setTicker
(String ticker)
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:32
Method
setTicker
(String ticker)
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:32
Method
setTicker
(String ticker)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:32
← previous
next →
401–500 of 544, ranked by callers