Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/apache/rocketmq-flink
/ functions
Functions
733 in github.com/apache/rocketmq-flink
⨍
Functions
733
◇
Types & classes
161
Method
read
(RowData row, int pos)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:270
Method
read
(RowData consumedRow, int pos)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:214
Method
read
(BytesMessage message)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:244
Method
releaseOutputForSplit
(String splitId)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:93
Method
report
(long delay)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/MetricUtils.java:96
Method
requestServiceDiscovery
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:262
Method
requiredOption
(ConfigOption<?> option)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigValidator.java:94
Method
requiredOptions
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactory.java:62
Method
requiredOptions
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:57
Method
restoreEnumerator
( SplitEnumeratorContext<RocketMQSourceSplit> enumContext, RocketMQSourceEnumState che
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:165
Method
resume
(Collection<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:229
Method
rocketMQSchema
( DeserializationSchema<T> valueDeserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:84
Method
rollback
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
Method
rollback
(SendCommittable sendCommittable)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:295
Method
run
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:88
Method
seek
(MessageQueue messageQueue, long offset)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:210
Method
seekCommittedOffset
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:235
Method
seekMaxOffset
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:289
Method
seekMinOffset
(MessageQueue messageQueue)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:267
Method
seekOffsetByTimestamp
( MessageQueue messageQueue, long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:311
Method
send
(Message message)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:159
Method
sendMessageInTransaction
(Message message)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:178
Method
serialize
(SendCommittable obj)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittableSerializer.java:38
Method
serializeKey
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueSerializationSchema.java:23
Method
serializeKeyAndValue
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueSerializationSchemaTest.java:29
Method
serializeValue
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueSerializationSchema.java:25
Method
setBodyOnlyDeserializer
( DeserializationSchema<OUT> deserializationSchema)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:134
Method
setBounded
(OffsetsSelector stoppingOffsetsSelector)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:122
Method
setConfig
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
Method
setConfigurationCanNotOverrideExistedKeysWithNewValue
()
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:27
Method
setEndpoints
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
Method
setGroupId
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
Method
setHeaderFields
(List<String> headerFields)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:520
Method
setMessageQueueSelector
( MessageQueueSelector messageQueueSelector)
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:95
Method
setProperties
(Map<String, String> properties)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:334
Method
setProperties
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
Method
setProperties
(Map<String, String> properties)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:525
Method
setPropertiesCanNotOverrideExistedKeysWithNewValueAndSupportTypeConversion
()
src/test/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilderTest.java:42
Method
setReSendInitAssign
(boolean reSendInitAssign)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceDetectEvent.java:31
Method
setRuntimeContext
(RuntimeContext runtimeContext)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataSink.java:44
Method
setTableSchema
(TableSchema tableSchema)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:299
Method
setUp
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSinkTest.java:45
Method
setUp
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceTest.java:58
Method
setUp
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:73
Method
snapshotState
(FunctionSnapshotContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSink.java:232
Method
snapshotState
(FunctionSnapshotContext context)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:563
Method
snapshotState
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceReader.java:104
Method
snapshotState
(long checkpointId)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:234
Method
start
()
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:137
Method
start
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:112
Method
start
()
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:167
Method
testAlterDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:246
Method
testAlterFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:236
Method
testAlterPartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:348
Method
testAlterPartitionColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:393
Method
testAlterPartitionStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:388
Method
testAlterTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:256
Method
testAlterTableColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:383
Method
testAlterTableStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:378
Method
testCall
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtilTest.java:34
Method
testClose
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:139
Method
testCreateCatalog
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:35
Method
testCreateDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:173
Method
testCreateFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:230
Method
testCreatePartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:336
Method
testCreateTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:204
Method
testDatabaseExists
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:167
Method
testDeserialize
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQRowDeserializationSchemaTest.java:26
Method
testDeserializeKeyAndValue
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchemaTest.java:35
Method
testDeserializeWithOldVersion
()
src/test/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializerTest.java:58
Method
testDeserializeWithOldVersion1
()
src/test/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializerTest.java:75
Method
testDropDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:178
Method
testDropFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:241
Method
testDropPartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:343
Method
testDropTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:209
Method
testEmitRecord
()
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:39
Method
testFactoryIdentifier
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:48
Method
testFunctionExists
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:225
Method
testGetDatabase
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:161
Method
testGetFactory
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:117
Method
testGetFunction
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:219
Method
testGetPartition
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:299
Method
testGetPartitionColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:371
Method
testGetPartitionStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:365
Method
testGetTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:190
Method
testGetTableColumnStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:359
Method
testGetTableStatistics
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:353
Method
testListDatabases
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:154
Method
testListFunctions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:214
Method
testListPartitions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:266
Method
testListPartitionsByFilter
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:293
Method
testListTables
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:183
Method
testListViews
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:251
Method
testOpen
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:123
Method
testOptionalOptions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:61
Method
testPartitionExists
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:321
Method
testRenameTable
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogTest.java:261
Method
testRequiredOptions
()
src/test/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactoryTest.java:54
Method
testRestartFromCheckpoint
()
src/test/java/org/apache/flink/connector/rocketmq/legacy/sourceFunction/RocketMQSourceFunctionTest.java:73
Method
testRocketMQDynamicTableSinkWithLegalOption
()
src/test/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSinkFactoryTest.java:60
← previous
next →
601–700 of 733, ranked by callers