Visitar URL original
fix: Change the return type of Heartbeat::getEstimatedLowWatermark to long by tengzhonger · Pull Request #1631 · googleapis/java-bigtable · GitHub
Skip to content
This repository was archived by the owner on Aug 30, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions google-cloud-bigtable/clirr-ignored-differences.xml
Original file line number Diff line number Diff line change
Expand Up @@ -86,4 +86,11 @@
<className>com/google/cloud/bigtable/gaxx/reframing/ReframingResponseObserver</className>
<to>com/google/api/gax/rpc/StateCheckingResponseObserver</to>
</difference>
<!-- change method return type is ok because Heartbeat is InternalApi -->
<difference>
<differenceType>7006</differenceType>
<className>com/google/cloud/bigtable/data/v2/models/Heartbeat</className>
<method>*getEstimatedLowWatermark*</method>
<to>long</to>
</difference>
</differences>
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
import com.google.api.core.InternalApi;
import com.google.auto.value.AutoValue;
import com.google.bigtable.v2.ReadChangeStreamResponse;
import com.google.protobuf.Timestamp;
import com.google.protobuf.util.Timestamps;
import java.io.Serializable;
import javax.annotation.Nonnull;

Expand All @@ -29,21 +29,20 @@ public abstract class Heartbeat implements ChangeStreamRecord, Serializable {
private static final long serialVersionUID = 7316215828353608504L;

private static Heartbeat create(
ChangeStreamContinuationToken changeStreamContinuationToken,
Timestamp estimatedLowWatermark) {
ChangeStreamContinuationToken changeStreamContinuationToken, long estimatedLowWatermark) {
return new AutoValue_Heartbeat(changeStreamContinuationToken, estimatedLowWatermark);
}

/** Wraps the protobuf {@link ReadChangeStreamResponse.Heartbeat}. */
static Heartbeat fromProto(@Nonnull ReadChangeStreamResponse.Heartbeat heartbeat) {
return create(
ChangeStreamContinuationToken.fromProto(heartbeat.getContinuationToken()),
heartbeat.getEstimatedLowWatermark());
Timestamps.toNanos(heartbeat.getEstimatedLowWatermark()));
}

@InternalApi("Intended for use by the BigtableIO in apache/beam only.")
public abstract ChangeStreamContinuationToken getChangeStreamContinuationToken();

@InternalApi("Intended for use by the BigtableIO in apache/beam only.")
public abstract Timestamp getEstimatedLowWatermark();
public abstract long getEstimatedLowWatermark();
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import com.google.cloud.bigtable.data.v2.models.Range.ByteStringRange;
import com.google.protobuf.ByteString;
import com.google.protobuf.Timestamp;
import com.google.protobuf.util.Timestamps;
import com.google.rpc.Status;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
Expand Down Expand Up @@ -119,7 +120,8 @@ public void heartbeatTest() {
.build();
Heartbeat actualHeartbeat = Heartbeat.fromProto(heartbeatProto);

assertThat(actualHeartbeat.getEstimatedLowWatermark()).isEqualTo(lowWatermark);
assertThat(actualHeartbeat.getEstimatedLowWatermark())
.isEqualTo(Timestamps.toNanos(lowWatermark));
assertThat(actualHeartbeat.getChangeStreamContinuationToken().getPartition())
.isEqualTo(ByteStringRange.create(rowRange.getStartKeyClosed(), rowRange.getEndKeyOpen()));
assertThat(actualHeartbeat.getChangeStreamContinuationToken().getToken()).isEqualTo(token);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import com.google.cloud.bigtable.gaxx.testing.FakeStreamingApi.ServerStreamingStashCallable;
import com.google.protobuf.ByteString;
import com.google.protobuf.Timestamp;
import com.google.protobuf.util.Timestamps;
import com.google.rpc.Status;
import java.util.Collections;
import java.util.List;
Expand Down Expand Up @@ -80,7 +81,7 @@ public void heartbeatTest() {
assertThat(heartbeat.getChangeStreamContinuationToken().getToken())
.isEqualTo(heartbeatProto.getContinuationToken().getToken());
assertThat(heartbeat.getEstimatedLowWatermark())
.isEqualTo(heartbeatProto.getEstimatedLowWatermark());
.isEqualTo(Timestamps.toNanos(heartbeatProto.getEstimatedLowWatermark()));
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,8 @@ public void test() throws Exception {
.build())
.setToken(heartbeat.getChangeStreamContinuationToken().getToken())
.build())
.setEstimatedLowWatermark(heartbeat.getEstimatedLowWatermark())
.setEstimatedLowWatermark(
Timestamps.fromNanos(heartbeat.getEstimatedLowWatermark()))
.build();
actualResults.add(
ReadChangeStreamTest.Result.newBuilder()
Expand Down