MCPcopy Create free account

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

Functions733 in github.com/apache/rocketmq-flink

↓ 2 callersMethodcalculateSplitAssignment
Calculate new split assignment according allocate strategy
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumerator.java:401
↓ 2 callersMethodclose
()
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:54
↓ 2 callersMethodcollect
(E record)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:72
↓ 2 callersMethoddatabaseExists
(String databaseName)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:155
↓ 2 callersMethodendTransaction
( final SendCommittable sendCommittable, final TransactionResult transactionResult)
src/main/java/org/apache/flink/connector/rocketmq/sink/InnerProducerImpl.java:227
↓ 2 callersMethodenqueueOffsetsCommitTask
( SplitFetcher<MessageView, RocketMQSourceSplit> splitFetcher, Map<MessageQueue, Long>
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:80
↓ 2 callersMethodgetCurrentOffsetTable
()
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:246
↓ 2 callersMethodgetMessageId
Get the unique message ID. @return the message ID
src/main/java/org/apache/flink/connector/rocketmq/source/reader/MessageView.java:31
↓ 2 callersMethodgetMessageOffset
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:117
↓ 2 callersMethodgetOffsetMsgId
()
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:97
↓ 2 callersMethodgetTag
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:37
↓ 2 callersMethodgetTag
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelector.java:67
↓ 2 callersMethodgetTopic
(T tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/DefaultTopicSelector.java:32
↓ 2 callersMethodgetTopic
(Map tuple)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/selector/SimpleTopicSelector.java:53
↓ 2 callersMethodgetVersion
()
src/main/java/org/apache/flink/connector/rocketmq/source/split/RocketMQPartitionSplitSerializer.java:36
↓ 2 callersMethodhandleSplitsChanges
(SplitsChange<RocketMQSourceSplit> splitsChange)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:156
↓ 2 callersMethodinitOffsetTableFromRestoredOffsets
(List<MessageQueue> messageQueues)
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:540
↓ 2 callersMethodinitOffsets
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 callersMethodisHeaderField
(int index)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:243
↓ 2 callersMethodlistPartitions
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:244
↓ 2 callersMethodminOffsets
List min offsets for the specified MessageQueues.
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/offset/OffsetsSelector.java:72
↓ 2 callersMethodnext
()
src/test/java/org/apache/flink/connector/rocketmq/sink/RocketMQSinkTest.java:70
↓ 2 callersMethodopen
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 callersMethodparseBoolean
(String s)
src/main/java/org/apache/flink/connector/rocketmq/source/util/StringSerializer.java:137
↓ 2 callersMethodparseDateString
(String dateString, String timeZone)
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQDynamicTableSourceFactory.java:184
↓ 2 callersMethodpause
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 callersMethodread
(BytesMessage message)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:404
↓ 2 callersMethodreadMetadata
(RowData consumedRow, WritableMetadata metadata)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:202
↓ 2 callersMethodsend
(MQProducer producer, String topic)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleProducer.java:59
↓ 2 callersMethodsetEndpoints
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 callersMethodsetGroupId
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 callersMethodsetMinOffsets
(OffsetsSelector offsetsSelector)
src/main/java/org/apache/flink/connector/rocketmq/source/RocketMQSourceBuilder.java:111
↓ 2 callersMethodsetMsgId
(String msgId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:93
↓ 2 callersMethodsetSplits
(List<RocketMQSourceSplit> splits)
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceInitAssignEvent.java:29
↓ 2 callersMethodsetStartFromGroupOffsets
consume from the group offsets those was stored in brokers.
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:461
↓ 2 callersMethodsetStartFromLatest
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 callersMethodsetTransactionId
(String transactionId)
src/main/java/org/apache/flink/connector/rocketmq/sink/committer/SendCommittable.java:109
↓ 2 callersMethodtableExists
(ObjectPath tablePath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:220
↓ 2 callersMethodwaitForMs
(long sleepMs)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RetryUtil.java:38
↓ 1 callersMethodallocate
( final Collection<RocketMQSourceSplit> mqAll, final int parallelism)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/ConsistentHashAllocateStrategy.java:41
↓ 1 callersMethodallocate
( final Collection<RocketMQSourceSplit> mqAll, final int parallelism)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/allocate/BroadcastAllocateStrategy.java:34
↓ 1 callersMethodalterDatabase
(String name, CatalogDatabase newDatabase, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:368
↓ 1 callersMethodalterFunction
( ObjectPath functionPath, CatalogFunction newFunction, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:355
↓ 1 callersMethodalterPartition
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogPartiti
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:419
↓ 1 callersMethodalterPartitionColumnStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogColumnS
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:481
↓ 1 callersMethodalterPartitionStatistics
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogTableSt
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:471
↓ 1 callersMethodalterTable
( ObjectPath tablePath, CatalogBaseTable newTable, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:380
↓ 1 callersMethodalterTableColumnStatistics
( ObjectPath tablePath, CatalogColumnStatistics columnStatistics, boolean
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:462
↓ 1 callersMethodalterTableStatistics
( ObjectPath tablePath, CatalogTableStatistics tableStatistics, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:455
↓ 1 callersMethodawaitTermination
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:361
↓ 1 callersMethodbuild
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:381
↓ 1 callersMethodbuild
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 callersMethodbuild
(RocketMQConfigValidator validator)
src/main/java/org/apache/flink/connector/rocketmq/common/config/RocketMQConfigBuilder.java:137
↓ 1 callersMethodbuildConsumerConfigs
Build Consumer Configs. @param props Properties @param consumer DefaultLitePullConsumer
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:150
↓ 1 callersMethodbuildExecutorService
(Configuration configuration)
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:134
↓ 1 callersMethodbuildProducerConfigs
Build Producer Configs. @param props Properties @param producer DefaultMQProducer
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQConfig.java:130
↓ 1 callersMethodbuilder
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 callersMethodbuilder
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 callersMethodcancel
()
src/main/java/org/apache/flink/connector/rocketmq/legacy/RocketMQSourceFunction.java:495
↓ 1 callersMethodclose
()
src/main/java/org/apache/flink/connector/rocketmq/source/InnerConsumerImpl.java:127
↓ 1 callersMethodcommit
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 callersMethodcommitOffset
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 callersMethodcommitOffsets
(Map<MessageQueue, Long> offsetsToCommit)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSourceFetcherManager.java:63
↓ 1 callersMethodcommonDeserialize
(byte[] value, ValueType type)
src/main/java/org/apache/flink/connector/rocketmq/source/util/ByteSerializer.java:61
↓ 1 callersMethodconvert
(RowData row)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQRowDataConverter.java:150
↓ 1 callersMethodconvertToRowTypeInfo
( DataType fieldsDataType, String[] fieldNames)
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:250
↓ 1 callersMethodcreateCatalog
(Context context)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:38
↓ 1 callersMethodcreateConverter
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:210
↓ 1 callersMethodcreateDatabase
(String name, CatalogDatabase database, boolean ignoreIfExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:307
↓ 1 callersMethodcreateFunction
( ObjectPath functionPath, CatalogFunction function, boolean ignoreIfExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:348
↓ 1 callersMethodcreateKeyValueDeserializationSchema
()
src/main/java/org/apache/flink/connector/rocketmq/source/table/RocketMQScanTableSource.java:212
↓ 1 callersMethodcreatePartition
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, CatalogPartiti
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:400
↓ 1 callersMethodcreateSink
()
src/main/java/org/apache/flink/connector/rocketmq/sink/table/RocketMQDynamicTableSink.java:206
↓ 1 callersMethodcreateTable
(ObjectPath tablePath, CatalogBaseTable table, boolean ignoreIfExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:319
↓ 1 callersMethodcreateTestTopic
(Set<String> brokerAddressSet)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:67
↓ 1 callersMethoddeleteSubscriptionGroup
(Set<String> brokerAddressSet)
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:91
↓ 1 callersMethoddeserialize
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 callersMethoddeserialize
(int version, byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:71
↓ 1 callersMethoddeserializeBytesMessage
(BytesMessage message, Collector<RowData> collector)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/deserializer/RowDeserializationSchema.java:198
↓ 1 callersMethoddeserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/KeyValueDeserializationSchema.java:24
↓ 1 callersMethoddeserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:105
↓ 1 callersMethoddeserializeKeyAndValue
(byte[] key, byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/SimpleKeyValueDeserializationSchema.java:49
↓ 1 callersMethoddeserializeMessageQueue
(byte[] serialized)
src/main/java/org/apache/flink/connector/rocketmq/source/enumerator/RocketMQSourceEnumStateSerializer.java:85
↓ 1 callersMethoddeserializeValue
(byte[] value)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/serialization/RowKeyValueDeserializationSchema.java:132
↓ 1 callersMethoddropDatabase
(String name, boolean ignoreIfNotExists, boolean cascade)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:313
↓ 1 callersMethoddropFunction
(ObjectPath functionPath, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:362
↓ 1 callersMethoddropPartition
( ObjectPath tablePath, CatalogPartitionSpec partitionSpec, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:412
↓ 1 callersMethoddropTable
(ObjectPath tablePath, boolean ignoreIfNotExists)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:325
↓ 1 callersMethodearliest
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 callersMethodemitRecord
( MessageView element, SourceOutput<T> output, RocketMQSourceSplitState splitState)
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitter.java:40
↓ 1 callersMethodemitWatermark
(Watermark watermark)
src/test/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQRecordEmitterTest.java:78
↓ 1 callersMethodextractTimestamp
(long timestamp)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/watermark/WaterMarkForAll.java:34
↓ 1 callersMethodfactoryIdentifier
()
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalogFactory.java:50
↓ 1 callersMethodfinishSplitAtRecord
( MessageQueue messageQueue, long currentOffset, RocketMQRecordsWithSplitI
src/main/java/org/apache/flink/connector/rocketmq/source/reader/RocketMQSplitReader.java:262
↓ 1 callersMethodflinkBodyOnlySchema
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 callersMethodfunctionExists
(ObjectPath functionPath)
src/main/java/org/apache/flink/connector/rocketmq/catalog/RocketMQCatalog.java:343
↓ 1 callersMethodgetAccessChannel
( Properties props, String key, AccessChannel defaultValue)
src/main/java/org/apache/flink/connector/rocketmq/legacy/common/util/RocketMQUtils.java:43
↓ 1 callersMethodgetAssignedMq
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceCheckEvent.java:31
↓ 1 callersMethodgetBroker
()
src/main/java/org/apache/flink/connector/rocketmq/common/event/SourceReportOffsetEvent.java:53
↓ 1 callersMethodgetBrokerAddress
()
src/test/java/org/apache/flink/connector/rocketmq/example/SimpleAdmin.java:59
← previousnext →101–200 of 733, ranked by callers