You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by ch...@apache.org on 2022/08/10 15:10:09 UTC
[beam] branch master updated (5d6f9f16dc6 -> fa9691fe2e9)
This is an automated email from the ASF dual-hosted git repository.
chamikara pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
from 5d6f9f16dc6 Improve exception when requested error tag does not exist (#22401) (#22405)
add fa9691fe2e9 Reimplement Pub/Sub Lite's I/O using UnboundedSource. (#22612)
No new revisions were added by this update.
Summary of changes:
.../beam/sdk/io/gcp/pubsub/TestPubsubSignal.java | 121 ++++++------
.../beam/sdk/io/gcp/pubsublite/PubsubLiteIO.java | 3 +
...ionPartitionProcessor.java => ApiServices.java} | 14 +-
.../gcp/pubsublite/internal/BlockingCommitter.java | 3 +-
...tReaderImpl.java => BlockingCommitterImpl.java} | 44 +++--
.../pubsublite/internal/CheckpointMarkImpl.java | 76 ++++++++
.../internal/LimitingTopicBacklogReader.java | 2 +-
.../internal/ManagedBacklogReaderFactoryImpl.java | 68 -------
...cklogReaderFactory.java => ManagedFactory.java} | 12 +-
.../pubsublite/internal/ManagedFactoryImpl.java | 60 ++++++
.../internal/PerSubscriptionPartitionSdf.java | 22 ++-
.../pubsublite/internal/SubscribeTransform.java | 85 ++++++---
.../pubsublite/internal/SubscriberAssembler.java | 59 ++++--
.../SubscriptionPartitionProcessorImpl.java | 42 +----
.../pubsublite/internal/TopicBacklogReader.java | 4 +-
.../pubsublite/internal/UnboundedReaderImpl.java | 144 +++++++++++++++
.../pubsublite/internal/UnboundedSourceImpl.java | 121 ++++++++++++
.../internal/BlockingCommmitterImplTest.java | 64 +++++++
.../internal/CheckpointMarkImplTest.java | 64 +++++++
.../internal/PerSubscriptionPartitionSdfTest.java | 32 ++--
.../SubscriptionPartitionProcessorImplTest.java | 38 +---
.../internal/UnboundedReaderImplTest.java | 202 +++++++++++++++++++++
22 files changed, 975 insertions(+), 305 deletions(-)
copy sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/{SubscriptionPartitionProcessor.java => ApiServices.java} (76%)
copy sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/{InitialOffsetReaderImpl.java => BlockingCommitterImpl.java} (51%)
create mode 100644 sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/CheckpointMarkImpl.java
delete mode 100644 sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/ManagedBacklogReaderFactoryImpl.java
rename sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/{ManagedBacklogReaderFactory.java => ManagedFactory.java} (70%)
create mode 100644 sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/ManagedFactoryImpl.java
create mode 100644 sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/UnboundedReaderImpl.java
create mode 100644 sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/UnboundedSourceImpl.java
create mode 100644 sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/BlockingCommmitterImplTest.java
create mode 100644 sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/CheckpointMarkImplTest.java
create mode 100644 sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/pubsublite/internal/UnboundedReaderImplTest.java