This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
from c97791057 [INLONG-5423][Sort] Revise invalid partition time when
dispatch messages (#5438)
add 67d89b384 [INLONG-5418][SDK] Sort SDK support seek topic or partitions
to a given timestamp (#5441)
No new revisions were added by this update.
Summary of changes:
.../inlong/sdk/sort/api/InLongTopicFetcher.java | 20 +++-
.../apache/inlong/sdk/sort/api/Interceptor.java | 9 +-
.../sdk/sort/api/{Interceptor.java => Seeker.java} | 31 +++---
.../api/{EmptyListener.java => SeekerFactory.java} | 32 +++----
.../apache/inlong/sdk/sort/entity/InLongTopic.java | 9 +-
.../sdk/sort/impl/QueryConsumeConfigImpl.java | 1 +
.../sort/impl/interceptor/MsgTimeInterceptor.java | 13 ++-
.../sdk/sort/impl/kafka/AckOffsetOnRebalance.java | 15 +--
.../sort/impl/kafka/InLongKafkaFetcherImpl.java | 5 +-
.../inlong/sdk/sort/impl/kafka/KafkaSeeker.java | 106 +++++++++++++++++++++
.../sort/impl/pulsar/InLongPulsarFetcherImpl.java | 32 ++++++-
.../inlong/sdk/sort/impl/pulsar/PulsarSeeker.java | 63 ++++++++++++
.../sdk/sort/impl/tube/InLongTubeFetcherImpl.java | 5 +-
.../impl/pulsar/InLongPulsarFetcherImplTest.java | 1 +
.../sort/standalone/sink/cls/ClsSinkContext.java | 18 ++--
15 files changed, 297 insertions(+), 63 deletions(-)
copy
inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/{Interceptor.java
=> Seeker.java} (64%)
copy
inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/{EmptyListener.java
=> SeekerFactory.java} (54%)
create mode 100644
inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/kafka/KafkaSeeker.java
create mode 100644
inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/impl/pulsar/PulsarSeeker.java