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
↓ 2 callers
Method
calculateSplitAssignment
Calculate new split assignment according allocate strategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:401
↓ 2 callers
Method
close
()
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:54
↓ 2 callers
Method
collect
(E record)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:72
↓ 2 callers
Method
databaseExists
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:155
↓ 2 callers
Method
endTransaction
( final SendCommittable sendCommittable, final TransactionResult transactionResult)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:227
↓ 2 callers
Method
enqueueOffsetsCommitTask
( SplitFetcher<MessageView, RocketMQSourceSplit> splitFetcher, Map<MessageQueue, Long>
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:80
↓ 2 callers
Method
getCurrentOffsetTable
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:246
↓ 2 callers
Method
getMessageId
Get the unique message ID. @return the message ID
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:31
↓ 2 callers
Method
getMessageOffset
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:117
↓ 2 callers
Method
getOffsetMsgId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:97
↓ 2 callers
Method
getTag
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:37
↓ 2 callers
Method
getTag
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelector.java:67
↓ 2 callers
Method
getTopic
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:32
↓ 2 callers
Method
getTopic
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelector.java:53
↓ 2 callers
Method
getVersion
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:36
↓ 2 callers
Method
handleSplitsChanges
(SplitsChange<RocketMQSourceSplit> splitsChange)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:156
↓ 2 callers
Method
initOffsetTableFromRestoredOffsets
(List<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:540
↓ 2 callers
Method
initOffsets
only flink job start with no state can init offsets from broker @param messageQueues @throws MQClientException
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:373
↓ 2 callers
Method
isHeaderField
(int index)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:243
↓ 2 callers
Method
listPartitions
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:244
↓ 2 callers
Method
minOffsets
List min offsets for the specified MessageQueues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:72
↓ 2 callers
Method
next
()
src/test/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkTest.java:70
↓ 2 callers
Method
open
Initialization method for the schema. It is called before the actual working methods {@link #serialize(Object, RocketMQSinkContext, Long)} and thus su
src/main/java/org/apache/flink/connector/rocketmq/sink/writer/serializer/RocketMQSerializationSchema.java:31
↓ 2 callers
Method
parseBoolean
(String s)
src/main/java/org/apache/flink/connector/rocketmq/source/util/StringSerializer.java:137
↓ 2 callers
Method
parseDateString
(String dateString, String timeZone)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:184
↓ 2 callers
Method
pause
Suspending message pulling from the message queues. @param messageQueues message queues that need to be suspended.
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:78
↓ 2 callers
Method
read
(BytesMessage message)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:404
↓ 2 callers
Method
readMetadata
(RowData consumedRow, WritableMetadata metadata)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:202
↓ 2 callers
Method
send
(MQProducer producer, String topic)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleProducer.java:59
↓ 2 callers
Method
setEndpoints
Configure the access point with which the SDK should communicate. @param endpoints address of service. @return the client configuration builder insta
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkBuilder.java:68
↓ 2 callers
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/sink/RocketMQSinkBuilder.java:78
↓ 2 callers
Method
setMinOffsets
(OffsetsSelector offsetsSelector)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:111
↓ 2 callers
Method
setMsgId
(String msgId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:93
↓ 2 callers
Method
setSplits
(List<RocketMQSourceSplit> splits)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceInitAssignEvent.java:29
↓ 2 callers
Method
setStartFromGroupOffsets
consume from the group offsets those was stored in brokers.
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:461
↓ 2 callers
Method
setStartFromLatest
consume from the max offset of each broker's queue at every restart with no state
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:448
↓ 2 callers
Method
setTransactionId
(String transactionId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:109
↓ 2 callers
Method
tableExists
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:220
↓ 2 callers
Method
waitForMs
(long sleepMs)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:38
↓ 1 callers
Method
allocate
( final Collection<RocketMQSourceSplit> mqAll, final int parallelism)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:41
↓ 1 callers
Method
allocate
( final Collection<RocketMQSourceSplit> mqAll, final int parallelism)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategy.java:34
↓ 1 callers
Method
alterDatabase
(String name, CatalogDatabase newDatabase, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:368
↓ 1 callers
Method
alterFunction
( ObjectPath functionPath, CatalogFunction newFunction, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:355
↓ 1 callers
Method
alterPartition
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogPartiti
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:419
↓ 1 callers
Method
alterPartitionColumnStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogColumnS
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:481
↓ 1 callers
Method
alterPartitionStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogTableSt
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:471
↓ 1 callers
Method
alterTable
( ObjectPath tablePath, CatalogBaseTable newTable, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:380
↓ 1 callers
Method
alterTableColumnStatistics
( ObjectPath tablePath, CatalogColumnStatistics columnStatistics, boolean
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:462
↓ 1 callers
Method
alterTableStatistics
( ObjectPath tablePath, CatalogTableStatistics tableStatistics, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:455
↓ 1 callers
Method
awaitTermination
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:361
↓ 1 callers
Method
build
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:381
↓ 1 callers
Method
build
Build the {@link RocketMQSource}. @return a RocketMQSource with the settings made for this builder.
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:183
↓ 1 callers
Method
build
(RocketMQConfigValidator validator)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:137
↓ 1 callers
Method
buildConsumerConfigs
Build Consumer Configs. @param props Properties @param consumer DefaultLitePullConsumer
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:150
↓ 1 callers
Method
buildExecutorService
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:134
↓ 1 callers
Method
buildProducerConfigs
Build Producer Configs. @param props Properties @param producer DefaultMQProducer
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:130
↓ 1 callers
Method
builder
Create a {@link RocketMQSinkBuilder} to construct a new {@link RocketMQSink}. @param <IN> type of incoming records @return {@link RocketMQSinkBuilder
src/main/java/org/apache/flink/connector/rocketmq/sink/RocketMQSink.java:54
↓ 1 callers
Method
builder
Get a RocketMQSourceBuilder to build a {@link RocketMQSourceBuilder}. @return a RocketMQ source builder.
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSource.java:95
↓ 1 callers
Method
cancel
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:495
↓ 1 callers
Method
close
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:127
↓ 1 callers
Method
commit
Commits the send operation identified by the specified SendCommittable object. @param sendCommittable the SendCommittable object identifying the send
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducer.java:73
↓ 1 callers
Method
commitOffset
Seek consumer group previously committed offset @param messageQueue rocketmq queue to locate single queue @return offset for message queue
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumer.java:136
↓ 1 callers
Method
commitOffsets
(Map<MessageQueue, Long> offsetsToCommit)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:63
↓ 1 callers
Method
commonDeserialize
(byte[] value, ValueType type)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:61
↓ 1 callers
Method
convert
(RowData row)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:150
↓ 1 callers
Method
convertToRowTypeInfo
( DataType fieldsDataType, String[] fieldNames)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:250
↓ 1 callers
Method
createCatalog
(Context context)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:38
↓ 1 callers
Method
createConverter
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:210
↓ 1 callers
Method
createDatabase
(String name, CatalogDatabase database, boolean ignoreIfExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:307
↓ 1 callers
Method
createFunction
( ObjectPath functionPath, CatalogFunction function, boolean ignoreIfExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:348
↓ 1 callers
Method
createKeyValueDeserializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:212
↓ 1 callers
Method
createPartition
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogPartiti
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:400
↓ 1 callers
Method
createSink
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:206
↓ 1 callers
Method
createTable
(ObjectPath tablePath, CatalogBaseTable table, boolean ignoreIfExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:319
↓ 1 callers
Method
createTestTopic
(Set<String> brokerAddressSet)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:67
↓ 1 callers
Method
deleteSubscriptionGroup
(Set<String> brokerAddressSet)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:91
↓ 1 callers
Method
deserialize
Deserializes the byte message. <p>Can output multiple records through the {@link Collector}. Note that number and size of the produced records should
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/QueryableSchema.java:55
↓ 1 callers
Method
deserialize
(int version, byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:71
↓ 1 callers
Method
deserializeBytesMessage
(BytesMessage message, Collector<RowData> collector)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:198
↓ 1 callers
Method
deserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueDeserializationSchema.java:24
↓ 1 callers
Method
deserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:105
↓ 1 callers
Method
deserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueDeserializationSchema.java:49
↓ 1 callers
Method
deserializeMessageQueue
(byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:85
↓ 1 callers
Method
deserializeValue
(byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:132
↓ 1 callers
Method
dropDatabase
(String name, boolean ignoreIfNotExists, boolean cascade)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:313
↓ 1 callers
Method
dropFunction
(ObjectPath functionPath, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:362
↓ 1 callers
Method
dropPartition
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:412
↓ 1 callers
Method
dropTable
(ObjectPath tablePath, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:325
↓ 1 callers
Method
earliest
Get an {@link OffsetsSelector} which initializes the offsets to the earliest available offsets of each partition. @return an {@link OffsetsSelector}
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:127
↓ 1 callers
Method
emitRecord
( MessageView element, SourceOutput<T> output, RocketMQSourceSplitState splitState)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:40
↓ 1 callers
Method
emitWatermark
(Watermark watermark)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:78
↓ 1 callers
Method
extractTimestamp
(long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkForAll.java:34
↓ 1 callers
Method
factoryIdentifier
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:50
↓ 1 callers
Method
finishSplitAtRecord
( MessageQueue messageQueue, long currentOffset, RocketMQRecordsWithSplitI
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:262
↓ 1 callers
Method
flinkBodyOnlySchema
Wraps a {@link DeserializationSchema} as the value deserialization schema. The other fields such as key, headers, timestamp are ignored. @param deser
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RocketMQDeserializationSchema.java:79
↓ 1 callers
Method
functionExists
(ObjectPath functionPath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:343
↓ 1 callers
Method
getAccessChannel
( Properties props, String key, AccessChannel defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:43
↓ 1 callers
Method
getAssignedMq
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceCheckEvent.java:31
↓ 1 callers
Method
getBroker
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:53
↓ 1 callers
Method
getBrokerAddress
()
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:59
← previous
next →
101–200 of 733, ranked by callers