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

↓ 1 callersMethodcreateDataGeneratorSource
Create a DataGeneratorSource with configurable rate from DataGen properties
java/JdbcSink/src/main/java/com/amazonaws/services/msf/JdbcSinkJob.java:67
↓ 1 callersMethodcreateDataGeneratorSource
(Properties dataGenProperties)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:86
↓ 1 callersMethodcreateDataGeneratorSource
(Properties dataGenProperties)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:60
↓ 1 callersMethodcreateDataStream
(StreamExecutionEnvironment env, Map<String, Properties> applicationProperties, Schema avroSchema)
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:109
↓ 1 callersMethodcreateKafkaSink
( Properties sinkProperties, Properties authProperties, KafkaRecordSeriali
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:51
↓ 1 callersMethodcreateKafkaSink
Create a KafkaSink @param kafkaProperties Properties from the "KafkaSink" property group @param recordSerializationSchema Record serializat
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:137
↓ 1 callersMethodcreateKafkaSink
(Properties outputProperties, KafkaRecordSerializationSchema<T> recordSerializationSchema)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:74
↓ 1 callersMethodcreateKafkaSink
Create the KafkaSink instance
java/KafkaConfigProviders/Kafka-SASL_SSL-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:78
↓ 1 callersMethodcreateKafkaSink
(Properties outputProperties, Tuple2<String, String> saslScramCredentials)
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:97
↓ 1 callersMethodcreateKafkaSource
( Properties sourceProperties, Properties authProperties, KafkaRecordDeser
java/AvroGlueSchemaRegistryKafka/consumer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:47
↓ 1 callersMethodcreateKafkaSource
(Properties inputProperties, final DeserializationSchema<T> valueDeserializationSchema)
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/KafkaStreamingJob.java:59
↓ 1 callersMethodcreateKafkaSource
(Properties applicationProperties)
java/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:61
↓ 1 callersMethodcreateKinesisFirehoseSink
( Properties sinkProperties)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/FirehoseStreamingJob.java:54
↓ 1 callersMethodcreateKinesisSink
(Properties outputProperties, final SerializationSchema<T> serializationSchema)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/StreamingJob.java:63
↓ 1 callersMethodcreateKinesisSink
(Properties outputProperties, final SerializationSchema<T> serializationSchema)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:53
↓ 1 callersMethodcreateKinesisSink
(Properties outputProperties, final SerializationSchema<T> serializationSchema)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/StreamingJob.java:59
↓ 1 callersMethodcreateKinesisSink
Create a Kinesis Sink @param outputProperties Properties from the "KinesisSink" property group @param serializationSchema Serialization schema
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:118
↓ 1 callersMethodcreateKinesisSink
(Properties outputProperties)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:59
↓ 1 callersMethodcreateKinesisSink
(Properties outputProperties, final SerializationSchema<T> serializationSchema)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/S3ParquetToKinesisJob.java:68
↓ 1 callersMethodcreateKinesisSink
(Properties kinesisClientProperties)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/CustomTypeInfoJob.java:58
↓ 1 callersMethodcreateKinesisSource
(Properties inputProperties, final KinesisDeserializationSchema<T> kinesisDeserializationSchema)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/StreamingJob.java:53
↓ 1 callersMethodcreateKinesisSource
(Properties inputProperties, final DeserializationSchema<T> deserializationSchema)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/StreamingJob.java:50
↓ 1 callersMethodcreateParquetS3Source
(Properties applicationProperties, final Class<T> clazz)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/S3ParquetToKinesisJob.java:52
↓ 1 callersMethodcreatePrometheusSink
(Properties prometheusSinkProperties)
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:71
↓ 1 callersMethodcreateSQSSink
( Properties sinkProperties, SqsSinkElementConverter<T> elementConverter)
java/SQSSink/src/main/java/com/amazonaws/services/msf/SQSStreamingJob.java:51
↓ 1 callersMethodcreateSink
(Properties outputProperties)
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/RecordCountJob.java:84
↓ 1 callersMethodcreateSink
(Properties outputProperties)
java/AsyncIO/src/main/java/com/amazonaws/services/msf/RetriesFlinkJob.java:101
↓ 1 callersMethodcreateSink
(Properties outputProperties)
java/GettingStarted/src/main/java/com/amazonaws/services/msf/BasicStreamingJob.java:52
↓ 1 callersMethodcreateSource
(Properties inputProperties)
java/GettingStarted/src/main/java/com/amazonaws/services/msf/BasicStreamingJob.java:45
↓ 1 callersMethodcreateTable
If Iceberg Table has not been previously created, we will create it using the Partition Fields specified in the Properties, as well as add a Sort Fiel
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:36
↓ 1 callersMethodcreateTable
If Iceberg Table has not been previously created, we will create it using the Partition Fields specified in the Properties, as well as add a Sort Fiel
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:46
↓ 1 callersMethodcreateTableStatement
(String sinkTableName)
java/Iceberg/IcebergSQLSink/src/main/java/com/amazonaws/services/msf/IcebergSQLSinkJob.java:75
↓ 1 callersMethodcreateUpsertJdbcSink
Create the JDBC Sink
java/JdbcSink/src/main/java/com/amazonaws/services/msf/JdbcSinkJob.java:89
↓ 1 callersMethoddeaggregate
Method to deaggregate a single Kinesis Record into a List of UserRecords @param inputRecord The Kinesis Record provided by AWS Lambda or Kinesis SDK
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/deaggregation/KinesisDeaggregatingDeserializationSchemaWrapper.java:81
↓ 1 callersFunctiondemo_flink_json
()
python/DatastreamKafkaConnector/datastream-kafka-connector-example.py:64
↓ 1 callersMethoddeserialize
(Record record, String s, String s1, Collector<ChangeEvent> collector)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/ChangeEventDeserializationSchema.java:16
↓ 1 callersMethodextractLongParameter
Extract a long parameter from Properties with validation and default fallback The parameter must be positive (> 0) @param properties The Properties o
java/JdbcSink/src/main/java/com/amazonaws/services/msf/ConfigurationHelper.java:106
↓ 1 callersMethodextractStringParameter
Extract an optional string parameter from Properties with default fallback @param properties The Properties object to extract from @param parameterNa
java/JdbcSink/src/main/java/com/amazonaws/services/msf/ConfigurationHelper.java:76
↓ 1 callersMethodfetchCredentialsFromSecretsManager
Fetch the secret from Secrets Manager. In the case of MSK SASL/SCRAM credentials, the secret is a JSON object with the keys `username` and `password`.
java/FetchSecrets/src/main/java/com/amazonaws/services/msf/FetchSecretsJob.java:61
↓ 1 callersFunctiongenerate
(stream_name, kinesis_client)
python/data-generator/stock.py:33
↓ 1 callersMethodgenerateRandomStockPrice
()
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/KplAggregatingProducer.java:132
↓ 1 callersFunctiongenerate_kafka_data
(topic_name, bootstrap_servers)
python/data-generator/stock_kafka.py:12
↓ 1 callersMethodgetDynamoDbStreamsSource
(Properties sourceProperties)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:47
↓ 1 callersMethodgetEventTime
()
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:14
↓ 1 callersMethodgetEventTime
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/StockPrice.java:20
↓ 1 callersMethodgetIncomingEventGeneratorSource
()
java/AsyncIO/src/main/java/com/amazonaws/services/msf/RetriesFlinkJob.java:92
↓ 1 callersMethodgetIncomingEventGeneratorSource
()
java/SideOutputs/src/main/java/com/amazonaws/services/msf/SideOutputsFlinkJob.java:109
↓ 1 callersMethodgetParquetS3Sink
(String s3UrlPath)
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:61
↓ 1 callersMethodgetParquetS3Sink
(String s3UrlPath)
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:61
↓ 1 callersMethodgetPartitionSpec
Generate the PartitionSpec, if when creating table, you want it to be partitioned. If you are doing Upserts in your Iceberg Table, your Equality Field
java/Iceberg/S3TableSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:52
↓ 1 callersMethodgetPartitionSpec
Generate the PartitionSpec, if when creating table, you want it to be partitioned. If you are doing Upserts in your Iceberg Table, your Equality Field
java/Iceberg/IcebergDataStreamSink/src/main/java/com/amazonaws/services/msf/iceberg/IcebergSinkBuilder.java:61
↓ 1 callersMethodgetPrice
()
java/KinesisSourceDeaggregation/kpl-producer/src/main/java/com/amazonaws/services/kds/producer/model/StockPrice.java:30
↓ 1 callersMethodgetPrice
()
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:38
↓ 1 callersMethodgetRoomId
()
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:19
↓ 1 callersMethodgetSpeedDataGeneratorSource
()
java/CustomMetrics/src/main/java/com/amazonaws/services/msf/RecordCountJob.java:75
↓ 1 callersMethodgetStockPriceDataGeneratorSource
()
java/Windowing/src/main/java/com/amazonaws/services/msf/windowing/WindowStreamingJob.java:100
↓ 1 callersMethodgetStockPriceDataGeneratorSource
()
java/S3Sink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:66
↓ 1 callersMethodgetStockPriceDataGeneratorSource
()
java/S3AvroSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:52
↓ 1 callersMethodgetStockPriceDataGeneratorSource
()
java/SQSSink/src/main/java/com/amazonaws/services/msf/SQSStreamingJob.java:42
↓ 1 callersMethodgetStockPriceDataGeneratorSource
()
java/S3ParquetSink/src/main/java/com/amazonaws/services/msf/StreamingJob.java:52
↓ 1 callersMethodgetStockPriceDataGeneratorSource
()
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/FirehoseStreamingJob.java:45
↓ 1 callersMethodgetSymbol
()
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:22
↓ 1 callersMethodgetTicker
()
java/KafkaConnectors/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:31
↓ 1 callersMethodgetTimestamp
()
java/JdbcSink/src/main/java/com/amazonaws/services/msf/domain/StockPrice.java:30
↓ 1 callersMethodgetTimestamp
()
java/PrometheusSink/src/main/java/com/amazonaws/services/msf/domain/TemperatureSample.java:35
↓ 1 callersFunctionget_application_properties
()
python/IcebergSink/main.py:87
↓ 1 callersFunctionget_application_properties
()
python/Windowing/main.py:76
↓ 1 callersFunctionget_application_properties
()
python/S3Sink/main.py:91
↓ 1 callersFunctionget_application_properties
()
python/FirehoseSink/main.py:76
↓ 1 callersFunctionget_application_properties
()
python/DatastreamKafkaConnector/datastream-kafka-connector-example.py:33
↓ 1 callersFunctionget_application_properties
()
python/PythonDependencies/main.py:85
↓ 1 callersFunctionget_application_properties
()
python/GettingStarted/main.py:78
↓ 1 callersFunctionget_application_properties
()
python/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders-DataStream/main.py:20
↓ 1 callersFunctionget_application_properties
()
python/UDF/main.py:80
↓ 1 callersFunctionget_application_properties
()
python/HudiSink/main.py:96
↓ 1 callersFunctionget_application_properties
()
python/PackagedPythonDependencies/main.py:103
↓ 1 callersMethodhasHDFSDelegationToken
Indicates whether the user has an HDFS delegation token.
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:74
↓ 1 callersMethodhasHDFSDelegationToken
Indicates whether the user has an HDFS delegation token.
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:74
↓ 1 callersMethodisKerberosSecurityEnabled
(UserGroupInformation ugi)
java/Iceberg/IcebergSQLSink/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:40
↓ 1 callersMethodisKerberosSecurityEnabled
(UserGroupInformation ugi)
java/S3ParquetSource/src/main/java/org/apache/flink/runtime/util/HadoopUtils.java:40
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:35
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/StreamingJob.java:33
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/S3AvroSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:33
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/KinesisConnectors/src/main/java/com/amazonaws/services/msf/StreamingJob.java:31
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/FlinkDataGenerator/src/main/java/com/amazonaws/services/msf/DataGeneratorJob.java:45
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/SQSSink/src/main/java/com/amazonaws/services/msf/SQSStreamingJob.java:25
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/DynamoDBStreamSource/src/main/java/com/amazonaws/services/msf/StreamingJob.java:28
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/S3ParquetSource/src/main/java/com/amazonaws/services/msf/S3ParquetToKinesisJob.java:33
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/KinesisFirehoseSink/src/main/java/com/amazonaws/services/msf/FirehoseStreamingJob.java:28
↓ 1 callersMethodisLocal
Detects whether the application is running locally or deployed in a cluster
java/KafkaConfigProviders/Kafka-SASL_SSL-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:54
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/KafkaConfigProviders/Kafka-mTLS-Keystore-Sql-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:38
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/KafkaConfigProviders/Kafka-mTLS-Keystore-ConfigProviders/src/main/java/com/amazonaws/services/msf/StreamingJob.java:40
↓ 1 callersMethodisLocal
(StreamExecutionEnvironment env)
java/Serialization/CustomTypeInfo/src/main/java/com/amazonaws/services/msf/CustomTypeInfoJob.java:40
↓ 1 callersMethodkafkaRecordDeserializationSchema
(Properties schemaRegistryProperties, Class<T> recordClazz)
java/AvroGlueSchemaRegistryKafka/consumer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:70
↓ 1 callersMethodkafkaRecordSerializationSchema
( Class<T> recordClazz, SerializableFunction<T, String> keyExtractor, Prop
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:69
↓ 1 callersMethodkinesisSink
Amazon Kinesis Flink Sink using AWS Glue Schema Registry @param <T> record type @param payloadAvroClass AVRO-generated class f
java/AvroGlueSchemaRegistryKinesis/src/main/java/com/amazonaws/services/msf/StreamingJob.java:106
↓ 1 callersMethodkinesisSource
Amazon Kinesis Flink Consumer using AWS Glue Schema Registry @param <T> record type @param payloadAvroClass AVRO-generated class for the
java/AvroGlueSchemaRegistryKinesis/src/main/java/com/amazonaws/services/msf/StreamingJob.java:69
↓ 1 callersMethodloadApplicationProperties
(StreamExecutionEnvironment env)
java/AvroGlueSchemaRegistryKafka/producer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:39
↓ 1 callersMethodloadApplicationProperties
(StreamExecutionEnvironment env)
java/AvroGlueSchemaRegistryKafka/consumer/src/main/java/com/amazonaws/services/msf/StreamingJob.java:35
↓ 1 callersMethodloadApplicationProperties
Load application properties from Amazon Managed Service for Apache Flink runtime or from a local resource, when the environment is local
java/KinesisSourceDeaggregation/flink-app/src/main/java/com/amazonaws/services/msf/StreamingJob.java:40
← previousnext →101–200 of 544, ranked by callers