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
StockPrice
()
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:13
Method
StockPrice
(String eventTime, String ticker, double price)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:8
Method
StockPrice
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:11
Method
StockPrice
()
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:11
Method
StockPrice
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:14
Method
StockPrice
()
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:9
Method
StockPrice
()
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:15
Method
StockPrice
()
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:11
Method
StockPrice
()
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:11
Method
StockPrice
()
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:11
Method
StockPrice
()
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:14
Method
StockPrice
()
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:8
Method
StockPriceUpsertQueryStatement
Create an UPSERT stock price query statement for a given table name. Note that, while the values are passed at runtime, the table name must be defined
java/JdbcSink/src/main/java/com/amazonaws/services/msf/jdbc/StockPriceUpsertQueryStatement.java:65
Method
TemperatureSample
(String roomId, String sensorId, long timestamp, double temperature)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:9
Method
VehicleEvent
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:76
Method
areKerberosCredentialsValid
( UserGroupInformation ugi, boolean useTicketCache)
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:47
Method
areKerberosCredentialsValid
( UserGroupInformation ugi, boolean useTicketCache)
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:47
Function
ask_bedrock_for_fun_fact
(a_number)
python/PythonDependencies/main.py:123
Method
asyncInvoke
(IncomingEvent incomingEvent, ResultFuture<ProcessedEvent> resultFuture)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessingFunction.java:53
Function
celsius_to_fahrenheit
(celsius)
python/UDF/main.py:104
Method
constructor
(scope: Construct, id: string, props?: cdk.StackProps)
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.ts:33
Method
constructor
(scope, id, props)
infrastructure/AutoScaling/cdk/lib/kda-autoscaling-stack.js:7
Method
createAccumulator
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:125
Method
createConverter
()
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/AvroSpecificRecordBulkFormat.java:34
Method
createReusedAvroRecord
()
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/AvroSpecificRecordBulkFormat.java:29
Method
createTypeInfo
(Type t, Map<String, TypeInformation<?>> genericParameters)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:32
Method
createTypeInfo
(Type t, Map<String, TypeInformation<?>> genericParameters)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:41
Method
createTypeInfo
(Type t, Map<String, TypeInformation<?>> genericParameters)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:50
Method
deserialize
(Record record, String stream, String shardId, Collector<T> output)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:37
Function
determinant
(element1, element2, element3, element4)
python/PackagedPythonDependencies/main.py:139
Method
equals
(Object o)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:47
Method
equals
(Object o)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:83
Method
generateRecord
()
java/Iceberg/S3TableSink/src/test/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunctionTest.java:12
Method
getDuration
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:136
Method
getElements
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:49
Method
getEventTime
()
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:15
Method
getEventTime
()
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:20
Method
getEventTime
()
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:15
Method
getEventTime
()
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:23
Method
getEventTime
()
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:20
Method
getEventTime
()
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:20
Method
getEventTime
()
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:20
Method
getEventType
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:20
Method
getFieldOne
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:37
Method
getFieldTwo
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:45
Method
getHadoopConfiguration
( org.apache.flink.configuration.Configuration flinkConfiguration)
python/IcebergSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:10
Method
getHadoopConfiguration
This method has been re-implemented to always return a org.apache.hadoop.conf.Configuration
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:35
Method
getHadoopConfiguration
This method has been re-implemented to always return a org.apache.hadoop.conf.Configuration
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:35
Method
getLocalDateTime
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:145
Method
getMapOfElements
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:57
Method
getMessage
()
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessedEvent.java:19
Method
getNewItem
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:36
Method
getOldItem
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:28
Method
getPartitionKey
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:16
Method
getPrice
()
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:31
Method
getPrice
()
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:36
Method
getPrice
()
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:31
Method
getPrice
()
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:35
Method
getPrice
()
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:39
Method
getPrice
()
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:36
Method
getPrice
()
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:36
Method
getPrice
()
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:36
Method
getProducedType
()
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:52
Method
getProducedType
()
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/AvroSpecificRecordBulkFormat.java:39
Method
getProducedType
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEventDeserializationSchema.java:41
Method
getResult
(Double accumulator)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:137
Method
getSetOfElements
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:65
Method
getSortKey
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/DdbTableItem.java:27
Method
getSymbol
()
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:27
Method
getTicker
()
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/model/StockPrice.java:23
Method
getTicker
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:20
Method
getTicker
()
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:28
Method
getTicker
()
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:23
Method
getTicker
()
java/SQSSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:28
Method
getTicker
()
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/StockPrice.java:28
Method
getTicker
()
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPrice.java:28
Method
getTimestamp
()
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:19
Method
getTimestamp
()
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:12
Method
getVolumes
()
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:43
Method
getZonedDateTime
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:128
Function
handler
(event, context)
infrastructure/AutoScaling/cdk/resources/scaling/scaling.py:58
Method
hashCode
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:95
Method
isAboveSpeedLimit
(SpeedRecord value)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/SpeedLimitFilter.java:7
Method
isLocal
(StreamExecutionEnvironment env)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:37
Method
isMaxHadoopVersion
Checks if the Hadoop dependency is at most the given version.
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:98
Method
isMaxHadoopVersion
Checks if the Hadoop dependency is at most the given version.
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:98
Method
isMinHadoopVersion
Checks if the Hadoop dependency is at least the given version.
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:87
Method
isMinHadoopVersion
Checks if the Hadoop dependency is at least the given version.
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:87
Method
main
(String[] args)
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:109
Method
main
(String[] args)
java/AvroGlueSchemaRegistryKafka/consumer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:80
Method
main
(String[] args)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/StreamingJob.java:73
Method
main
(String[] args)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:47
Method
main
(String[] args)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:56
Method
main
(String[] args)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:79
Method
main
(String[] args)
java/S3Sink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:75
Method
main
(String[] args)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/StreamingJob.java:69
Method
main
(String[] args)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:146
Method
main
(String[] args)
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:74
Method
main
(String[] args)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/RecordCountJob.java:45
Method
main
(String[] args)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:114
← previous
next →
301–400 of 544, ranked by callers