MCPcopy Create free account

hub / github.com/apache/rocketmq-flink / functions

Functions733 in github.com/apache/rocketmq-flink

Methodread
(RowData row, int pos)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:270
Methodread
(RowData consumedRow, int pos)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:214
Methodread
(BytesMessage message)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:244
MethodreleaseOutputForSplit
(String splitId)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:93
Methodreport
(long delay)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:96
MethodrequestServiceDiscovery
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:262
MethodrequiredOption
(ConfigOption<?> option)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:94
MethodrequiredOptions
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:62
MethodrequiredOptions
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:57
MethodrestoreEnumerator
( SplitEnumeratorContext<RocketMQSourceSplit> enumContext, RocketMQSourceEnumState che
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:165
Methodresume
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:229
MethodrocketMQSchema
( DeserializationSchema<T> valueDeserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:84
Methodrollback
Rolls back the send operation identified by the specified SendCommittable object. @param sendCommittable the SendCommittable object identifying the s
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:83
Methodrollback
(SendCommittable sendCommittable)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:295
Methodrun
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:88
Methodseek
(MessageQueue messageQueue, long offset)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:210
MethodseekCommittedOffset
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:235
MethodseekMaxOffset
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:289
MethodseekMinOffset
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:267
MethodseekOffsetByTimestamp
( MessageQueue messageQueue, long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:311
Methodsend
(Message message)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:159
MethodsendMessageInTransaction
(Message message)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:178
Methodserialize
(SendCommittable obj)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittableSerializer.java:38
MethodserializeKey
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueSerializationSchema.java:23
MethodserializeKeyAndValue
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchemaTest.java:29
MethodserializeValue
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueSerializationSchema.java:25
MethodsetBodyOnlyDeserializer
( DeserializationSchema<OUT> deserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:134
MethodsetBounded
(OffsetsSelector stoppingOffsetsSelector)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:122
MethodsetConfig
Set an arbitrary property for the RocketMQ source. The valid keys can be found in {@link RocketMQSourceOptions}. Make sure the option could be set onl
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:149
MethodsetConfigurationCanNotOverrideExistedKeysWithNewValue
()
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:27
MethodsetEndpoints
Configure the access point with which the SDK should communicate. @param endpoints address of service. @return this RocketMQSourceBuilder.
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:70
MethodsetGroupId
Sets the consumer group id of the RocketMQSource. @param groupId the group id of the RocketMQSource. @return this RocketMQSourceBuilder.
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:81
MethodsetHeaderFields
(List<String> headerFields)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:520
MethodsetMessageQueueSelector
( MessageQueueSelector messageQueueSelector)
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:95
MethodsetProperties
(Map<String, String> properties)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:334
MethodsetProperties
Set arbitrary properties for the RocketMQ source. This method is mainly used for future flink SQL binding. @param properties the config properties to
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:173
MethodsetProperties
(Map<String, String> properties)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:525
MethodsetPropertiesCanNotOverrideExistedKeysWithNewValueAndSupportTypeConversion
()
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:42
MethodsetReSendInitAssign
(boolean reSendInitAssign)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceDetectEvent.java:31
MethodsetRuntimeContext
(RuntimeContext runtimeContext)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:44
MethodsetTableSchema
(TableSchema tableSchema)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:299
MethodsetUp
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSinkTest.java:45
MethodsetUp
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceTest.java:58
MethodsetUp
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:73
MethodsnapshotState
(FunctionSnapshotContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:232
MethodsnapshotState
(FunctionSnapshotContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:563
MethodsnapshotState
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:104
MethodsnapshotState
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:234
Methodstart
()
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:137
Methodstart
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:112
Methodstart
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:167
MethodtestAlterDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:246
MethodtestAlterFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:236
MethodtestAlterPartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:348
MethodtestAlterPartitionColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:393
MethodtestAlterPartitionStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:388
MethodtestAlterTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:256
MethodtestAlterTableColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:383
MethodtestAlterTableStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:378
MethodtestCall
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtilTest.java:34
MethodtestClose
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:139
MethodtestCreateCatalog
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:35
MethodtestCreateDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:173
MethodtestCreateFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:230
MethodtestCreatePartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:336
MethodtestCreateTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:204
MethodtestDatabaseExists
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:167
MethodtestDeserialize
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchemaTest.java:26
MethodtestDeserializeKeyAndValue
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchemaTest.java:35
MethodtestDeserializeWithOldVersion
()
src/test/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializerTest.java:58
MethodtestDeserializeWithOldVersion1
()
src/test/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializerTest.java:75
MethodtestDropDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:178
MethodtestDropFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:241
MethodtestDropPartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:343
MethodtestDropTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:209
MethodtestEmitRecord
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:39
MethodtestFactoryIdentifier
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:48
MethodtestFunctionExists
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:225
MethodtestGetDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:161
MethodtestGetFactory
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:117
MethodtestGetFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:219
MethodtestGetPartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:299
MethodtestGetPartitionColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:371
MethodtestGetPartitionStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:365
MethodtestGetTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:190
MethodtestGetTableColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:359
MethodtestGetTableStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:353
MethodtestListDatabases
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:154
MethodtestListFunctions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:214
MethodtestListPartitions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:266
MethodtestListPartitionsByFilter
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:293
MethodtestListTables
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:183
MethodtestListViews
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:251
MethodtestOpen
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:123
MethodtestOptionalOptions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:61
MethodtestPartitionExists
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:321
MethodtestRenameTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:261
MethodtestRequiredOptions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:54
MethodtestRestartFromCheckpoint
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/sourceFunction/RocketMQSourceFunctionTest.java:73
MethodtestRocketMQDynamicTableSinkWithLegalOption
()
src/test/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactoryTest.java:60
← previousnext →601–700 of 733, ranked by callers