MCPcopy Create free account

hub / github.com/aws-samples/amazon-managed-service-for-apache-flink-examples / functions

Functions544 in github.com/aws-samples/amazon-managed-service-for-apache-flink-examples

↓ 10 callersMethodgetTicker
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:30
↓ 8 callersMethodequals
(Object o)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:46
↓ 8 callersFunctionproperty_map
(props, property_group_id)
python/Windowing/main.py:87
↓ 7 callersMethodgetPrice
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:38
↓ 6 callersMethodhashCode
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:56
↓ 5 callersMethodgetEventTime
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:22
↓ 5 callersMethodgetSensorId
()
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:27
↓ 5 callersMethodgetTemperature
()
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:43
↓ 5 callersMethodhashCode
()
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:56
↓ 5 callersMethodtoString
()
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:51
↓ 4 callersMethodcreateSink
(Properties outputProperties)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:109
↓ 4 callersMethodgetAddress
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:180
↓ 4 callersMethodgetMinPrice
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:36
↓ 4 callersMethodgetName
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:172
↓ 4 callersMethodgetTicker
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:28
↓ 4 callersMethodgetTimeWindow
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/AggregatedStockPrice.java:28
↓ 4 callersMethodisLocal
(StreamExecutionEnvironment env)
java/FlinkCDC/FlinkCDCSQLServerSource/src/main/java/com/amazonaws/services/msf/FlinkCDCSqlServer2JdbcJob.java:25
↓ 4 callersMethodmap
(Long aLong)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPriceGeneratorFunction.java:12
↓ 4 callersFunctionproperty_map
(props, property_group_id)
python/PythonDependencies/main.py:95
↓ 4 callersMethodserializeDeserializeNoKryo
Round-trip serialization-deserialization of a record, using Flink serializers, disabling "generic types", i.e. disabling Kryo fallback. The serializat
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/FlinkSerializationTestUtils.java:66
↓ 4 callersMethodtoString
()
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:10
↓ 4 callersFunctionupdate_parallelism
(context, desiredCapacity, resourceName, appVersionId, currentParallelismPerKPU)
infrastructure/AutoScaling/cdk/resources/scaling/scaling.py:12
↓ 3 callersMethodequals
(Object o)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:128
↓ 3 callersMethodextractIntParameter
Extract an integer parameter from Properties with validation and default fallback The parameter must be positive (> 0) @param properties The Properti
java/JdbcSink/src/main/java/com/amazonaws/services/msf/ConfigurationHelper.java:93
↓ 3 callersMethodextractRequiredStringParameter
Extract a required string parameter from Properties Throws IllegalArgumentException if the parameter is missing or empty @param properties The Proper
java/JdbcSink/src/main/java/com/amazonaws/services/msf/ConfigurationHelper.java:60
↓ 3 callersMethodgetId
()
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:28
↓ 3 callersMethodgetMessage
()
java/AsyncIO/src/main/java/com/amazonaws/services/msf/IncomingEvent.java:24
↓ 3 callersMethodgetPrice
()
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:32
↓ 3 callersMethodgetSensorData
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:102
↓ 3 callersMethodgetSymbol
()
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:16
↓ 3 callersMethodgetWarningLights
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:110
↓ 3 callersMethodhashCode
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:140
↓ 3 callersMethodmap
(Long index)
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/domain/TemperatureSampleGenerator.java:14
↓ 3 callersMethodmap
(SpeedRecord value)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/MetricEmittingMapperFunction.java:38
↓ 3 callersMethodmap
(Long aLong)
java/SideOutputs/src/main/java/com/amazonaws/services/msf/IncomingEventDataGeneratorFunction.java:10
↓ 3 callersFunctionproperty_map
(props, property_group_id)
python/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders-DataStream/main.py:29
↓ 3 callersMethodsetEventType
(Type eventType)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:24
↓ 3 callersMethodtoString
()
java/S3Sink/src/main/java/com/amazonaws/services/msf/StockPrice.java:44
↓ 2 callersMethodadd
(StockPrice value, Double accumulator)
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:130
↓ 2 callersMethodcreateSink
(Properties outputProperties)
java/SideOutputs/src/main/java/com/amazonaws/services/msf/SideOutputsFlinkJob.java:118
↓ 2 callersMethodequals
(Object o)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:46
↓ 2 callersMethodextractNumericParameter
Generic method to extract and parse numeric parameters from Properties All parameters are validated to be positive (> 0) @param properties The Proper
java/JdbcSink/src/main/java/com/amazonaws/services/msf/ConfigurationHelper.java:23
↓ 2 callersMethodforAvroSchema
(Schema avroSchema)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/AvroGenericRecordToRowDataMapper.java:25
↓ 2 callersMethodgetAverageSensorData
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:67
↓ 2 callersMethodgetCountDistinctWarnings
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:75
↓ 2 callersMethodgetMajorMinorBundledHadoopVersion
()
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:106
↓ 2 callersMethodgetMajorMinorBundledHadoopVersion
()
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:106
↓ 2 callersMethodgetPostCode
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:201
↓ 2 callersMethodgetPrice
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:36
↓ 2 callersMethodgetStreet
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:193
↓ 2 callersMethodgetTicker
()
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:22
↓ 2 callersMethodgetTimestamp
()
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:24
↓ 2 callersMethodgetTimestamp
()
java/Serialization/CustomTypeInfo/src/test/java/com/amazonaws/services/msf/domain/MoreKryoSerializationExamplesTest.java:91
↓ 2 callersMethodgetTimestamp
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:59
↓ 2 callersMethodgetTimestamp
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:94
↓ 2 callersMethodgetValue
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/MetricEmitterNoOpMap.java:35
↓ 2 callersMethodgetVin
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/AggregateVehicleEvent.java:51
↓ 2 callersMethodgetVin
()
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/domain/VehicleEvent.java:86
↓ 2 callersFunctionget_data
()
python/data-generator/stock.py:26
↓ 2 callersMethodhashCode
()
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:57
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/AvroGlueSchemaRegistryKafka/consumer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:31
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/S3Sink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:34
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:33
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:34
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/Iceberg/IcebergDataStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:41
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:38
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:32
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:38
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:33
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/JdbcSinkJob.java:45
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:38
↓ 2 callersMethodisLocal
(StreamExecutionEnvironment env)
java/AvroGlueSchemaRegistryKinesis/src/main/java/com/amazonaws/services/msf/StreamingJob.java:40
↓ 2 callersMethodloadSchema
Load the AVRO Schema from the resources folder
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/avro/AvroSchemaUtils.java:17
↓ 2 callersMethodmap
(T record)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/JsonConverter.java:29
↓ 2 callersMethodmap
(Long value)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPriceGeneratorFunction.java:46
↓ 2 callersMethodmergeProperties
(Properties properties, Properties authProperties)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:83
↓ 2 callersMethodprocess
(String vin, ProcessWindowFunction<VehicleEvent, AggregateVehicleEvent, String, TimeWindow>.Context context, I
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/aggregation/AggregateVehicleEventsWindowFunction.java:22
↓ 2 callersFunctionproperty_map
(props, property_group_id)
python/FirehoseSink/main.py:87
↓ 2 callersFunctionproperty_map
(props, property_group_id)
python/DatastreamKafkaConnector/datastream-kafka-connector-example.py:43
↓ 2 callersFunctionproperty_map
(props, property_group_id)
python/GettingStarted/main.py:88
↓ 2 callersFunctionproperty_map
(props, property_group_id)
python/UDF/main.py:91
↓ 2 callersFunctionproperty_map
(props, property_group_id)
python/PackagedPythonDependencies/main.py:114
↓ 2 callersMethodsetNewItem
(DdbTableItem newItem)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:40
↓ 2 callersMethodsetOldItem
(DdbTableItem oldItem)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEvent.java:32
↓ 2 callersMethodsetTimestamp
(Instant timestamp)
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:34
↓ 2 callersMethodtoString
()
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:61
↓ 1 callersMethodS3Sink
@param s3SinkPath: Path to which application will write data to @return FileSink
java/S3Sink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:58
↓ 1 callersMethodaccept
Replace the positional parameters in the prepared statement. The implementation of this method depends on the SQL statement which, in turn, is specifi
java/JdbcSink/src/main/java/com/amazonaws/services/msf/jdbc/StockPriceUpsertQueryStatement.java:43
↓ 1 callersMethodclose
()
java/AsyncIO/src/main/java/com/amazonaws/services/msf/ProcessingFunction.java:47
↓ 1 callersMethodconfigureConnectorPropsWithConfigProviders
(KafkaSourceBuilder<String> builder, Properties appProperties)
java/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:78
↓ 1 callersMethodconvertType
(List<T> inputRecords)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:85
↓ 1 callersMethodcreateAvroFileSource
(Properties sourceProperties, Class<T> avroRecordClass, Schema avroSchema)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:63
↓ 1 callersMethodcreateBuilder
S3 Table Sink Builder - It is unique in that it leverages the S3 Table Catalog as opposed to an external catalog like Glue or Hive. catalog-impl = sof
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:64
↓ 1 callersMethodcreateBuilder
(Properties icebergProperties, DataStream<GenericRecord> dataStream, org.apache.avro.Schema avroSchema)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:70
↓ 1 callersMethodcreateCatalogStatement
(String s3BucketPrefix)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:67
↓ 1 callersMethodcreateDataGenerator
(Properties dataGeneratorProperties)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:56
↓ 1 callersMethodcreateDataGenerator
(Properties generatorProperties, Schema avroSchema)
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:58
↓ 1 callersMethodcreateDataGenerator
(Properties generatorProperties, Schema avroSchema)
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:53
↓ 1 callersMethodcreateDataGeneratorSource
Create a DataGeneratorSource with configurable rate from DataGen properties @param dataGenProperties Properties from the "DataGen" property group @pa
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:73
↓ 1 callersMethodcreateDataGeneratorSource
Create the source instance. For simplicity, we use a DataGeneratorSource to generate random strings. In a real application, this would be a connector
java/KafkaConfigProviders/Kafka-SASL_SSL-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:140
next →1–100 of 544, ranked by callers