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
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:44
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:40
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/S3Sink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:41
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/StreamingJob.java:38
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:52
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:40
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/RecordCountJob.java:33
↓ 1 callers
Method
loadApplicationProperties
(StreamExecutionEnvironment env)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:44
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/Iceberg/IcebergDataStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:48
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:45
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:39
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:46
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/SQSSink/src/main/java/com/amazonaws/services/msf/SQSStreamingJob.java:32
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/GettingStartedTable/src/main/java/com/amazonaws/services/msf/BasicTableJob.java:34
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:35
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:40
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/S3ParquetToKinesisJob.java:40
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/FlinkCDC/FlinkCDCSQLServerSource/src/main/java/com/amazonaws/services/msf/FlinkCDCSqlServer2JdbcJob.java:32
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/FirehoseStreamingJob.java:35
↓ 1 callers
Method
loadApplicationProperties
(StreamExecutionEnvironment env)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/RetriesFlinkJob.java:34
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/GettingStarted/src/main/java/com/amazonaws/services/msf/BasicStreamingJob.java:33
↓ 1 callers
Method
loadApplicationProperties
Get configuration properties from Amazon Managed Service for Apache Flink runtime properties or from a local resource when running locally
java/KafkaConfigProviders/Kafka-SASL_SSL-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:62
↓ 1 callers
Method
loadApplicationProperties
Get configuration properties from Amazon Managed Service for Apache Flink runtime properties or from a local resource when running locally
java/KafkaConfigProviders/Kafka-mTLS-Keystore-Sql-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:46
↓ 1 callers
Method
loadApplicationProperties
Get configuration properties from Amazon Managed Service for Apache Flink runtime properties or from a local resource when running locally
java/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:48
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/SideOutputs/src/main/java/com/amazonaws/services/msf/SideOutputsFlinkJob.java:46
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/JdbcSink/src/main/java/com/amazonaws/services/msf/JdbcSinkJob.java:52
↓ 1 callers
Method
loadApplicationProperties
(StreamExecutionEnvironment env)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:42
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:47
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/AvroGlueSchemaRegistryKinesis/src/main/java/com/amazonaws/services/msf/StreamingJob.java:47
↓ 1 callers
Method
loadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/CustomTypeInfoJob.java:47
↓ 1 callers
Method
loadSchema
Load the AVRO Schema from the resources folder
java/Iceberg/IcebergDataStreamSource/src/main/java/com/amazonaws/services/msf/avro/AvroSchemaUtils.java:17
↓ 1 callers
Method
loadSchema
Load the AVRO Schema from the resources folder
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/avro/AvroSchemaUtils.java:17
↓ 1 callers
Function
main
()
python/IcebergSink/main.py:102
↓ 1 callers
Function
main
()
python/Windowing/main.py:93
↓ 1 callers
Function
main
()
python/S3Sink/main.py:108
↓ 1 callers
Function
main
()
python/FirehoseSink/main.py:93
↓ 1 callers
Function
main
()
python/PythonDependencies/main.py:162
↓ 1 callers
Function
main
()
python/GettingStarted/main.py:96
↓ 1 callers
Function
main
()
python/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders-DataStream/main.py:34
↓ 1 callers
Function
main
()
python/UDF/main.py:118
↓ 1 callers
Function
main
()
python/HudiSink/main.py:111
↓ 1 callers
Function
main
()
python/PackagedPythonDependencies/main.py:151
↓ 1 callers
Method
map
(self, value)
python/DatastreamKafkaConnector/datastream-kafka-connector-example.py:55
↓ 1 callers
Method
map
(Long aLong)
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:12
↓ 1 callers
Method
map
(T record)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/MetricEmitterNoOpMap.java:45
↓ 1 callers
Method
map
Generates a random trade
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunction.java:30
↓ 1 callers
Method
map
(Long aLong)
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/datagen/StockPriceGeneratorFunction.java:12
↓ 1 callers
Method
map
(T record)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:29
↓ 1 callers
Method
map
(Long aLong)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/StockPriceGeneratorFunction.java:12
↓ 1 callers
Method
map
(Long aLong)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEventDataGeneratorFunction.java:9
↓ 1 callers
Method
map
(Long value)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPriceGenerator.java:13
↓ 1 callers
Method
map
(Long value)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/datagen/TemperatureSampleGeneratorFunction.java:13
↓ 1 callers
Method
merge
(Double a, Double b)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:142
↓ 1 callers
Method
open
(OpenContext openContext)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:24
↓ 1 callers
Method
open
(OpenContext openContext)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/MetricEmitterNoOpMap.java:19
↓ 1 callers
Method
produce
()
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:97
↓ 1 callers
Function
property_map
(props, property_group_id)
python/IcebergSink/main.py:97
↓ 1 callers
Function
property_map
(props, property_group_id)
python/S3Sink/main.py:102
↓ 1 callers
Function
property_map
(props, property_group_id)
python/HudiSink/main.py:106
↓ 1 callers
Function
response_function
(status_code, response_body)
infrastructure/AutoScaling/cdk/resources/scaling/scaling.py:45
↓ 1 callers
Method
schemaRegistryConf
Glue Schema Registry SerDe configuration @param schemaRegistryRegion Glue Schema Registry region @param schemaRegistryName Glue Schema Registry nam
java/AvroGlueSchemaRegistryKinesis/src/main/java/com/amazonaws/services/msf/StreamingJob.java:135
↓ 1 callers
Method
sendStockPriceToKinesis
(StockPrice stockPrice, String streamName, KinesisProducer producer,
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:140
↓ 1 callers
Method
serializeDeserialize
Round-trip serialization-deserialization of a record, using Flink serializers. Returns the deserialized record por fails, if Flink cannot serialize/de
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/FlinkSerializationTestUtils.java:33
↓ 1 callers
Method
setAddress
(Address address)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:184
↓ 1 callers
Method
setEventTime
(String eventTime)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:26
↓ 1 callers
Method
setId
(String id)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:32
↓ 1 callers
Method
setName
(String name)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:176
↓ 1 callers
Method
setPostCode
(String postCode)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:205
↓ 1 callers
Method
setPrice
(float price)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:42
↓ 1 callers
Method
setSensorId
(String sensorId)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:31
↓ 1 callers
Method
setStreet
(String street)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:197
↓ 1 callers
Method
setTemperature
(double temperature)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:47
↓ 1 callers
Method
setTicker
(String ticker)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:34
↓ 1 callers
Method
setTimestamp
(Instant timestamp)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:16
↓ 1 callers
Method
setTimestamp
(Instant timestamp)
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:95
↓ 1 callers
Method
setupIcebergProperties
(Properties icebergProperties)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:84
↓ 1 callers
Method
toByteArray
(ByteBuffer buffer)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:57
↓ 1 callers
Method
toString
()
java/SideOutputs/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:7
↓ 1 callers
Method
toString
()
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:40
↓ 1 callers
Method
toString
()
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:51
↓ 1 callers
Method
validateURI
(String uri)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:38
Method
AggregateVehicleEvent
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:41
Method
AggregatedStockPrice
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:11
Method
AvroGenericRecordToRowDataMapper
(RowType rowType)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/AvroGenericRecordToRowDataMapper.java:17
Method
AvroGenericStockTradeGeneratorFunction
(Schema avroSchema)
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunction.java:23
Method
AvroGenericStockTradeGeneratorFunction
(Schema avroSchema)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/datagen/AvroGenericStockTradeGeneratorFunction.java:23
Method
AvroSpecificRecordBulkFormat
(Class<O> recordClass, org.apache.avro.Schema avroSchema)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/AvroSpecificRecordBulkFormat.java:23
Method
IcebergConfig
(String s3BucketPrefix, String glueDatabase, String glueTable)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:107
Method
IncomingEvent
(String message)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:18
Method
IncomingEvent
(String message)
java/SideOutputs/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:14
Method
JsonConverter
(Schema avroSchema)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:20
Method
JsonConverter
(Schema avroSchema)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:20
Method
KinesisDeaggregatingDeserializationSchemaWrapper
(DeserializationSchema<T> deserializationSchema)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:28
Method
KplAggregatingProducer
(String streamName, String streamRegion, long sleepTimeBetweenRecords)
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:91
Method
MetricEmittingMapperFunction
(final String customMetricName)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/MetricEmittingMapperFunction.java:18
Method
ProcessedEvent
(String message)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessedEvent.java:14
Method
ProcessingFunction
(String apiUrl, String apiKey)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessingFunction.java:27
Method
ProcessingOutcome
(final String text)
java/SideOutputs/src/main/java/com/amazonaws/services/msf/ProcessingOutcome.java:9
Method
SpeedRecord
(String id, double speed)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/SpeedRecord.java:8
Method
Stock
()
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/Stock.java:13
← previous
next →
201–300 of 544, ranked by callers