Skip to content
Open
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
2 changes: 1 addition & 1 deletion sdk/cosmos/azure-cosmos-kafka-connect/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#### Breaking Changes

#### Bugs Fixed
* Fixed source connector data loss after partition splits when a stale parent continuation contained divergent child LSNs. - See [PR 50031](https://github.com/Azure/azure-sdk-for-java/pull/50031)

#### Other Changes

Expand Down Expand Up @@ -132,4 +133,3 @@
* Added `ServicePrincipal` support - See [PR 39490](https://github.com/Azure/azure-sdk-for-java/pull/39490)
* Added `ItemPatch support` in sink connector - See [PR 39558](https://github.com/Azure/azure-sdk-for-java/pull/39558)
* Added support to use CosmosDB container for tracking metadata - See [PR 39634](https://github.com/Azure/azure-sdk-for-java/pull/39634)

1 change: 1 addition & 0 deletions sdk/cosmos/azure-cosmos-kafka-connect/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ Licensed under the MIT License.
--add-exports com.azure.cosmos/com.azure.cosmos.implementation.changefeed.common=com.azure.cosmos.kafka.connect
--add-exports com.azure.cosmos/com.azure.cosmos.implementation.feedranges=com.azure.cosmos.kafka.connect
--add-exports com.azure.cosmos/com.azure.cosmos.implementation.query=com.azure.cosmos.kafka.connect
--add-exports com.azure.cosmos/com.azure.cosmos.implementation.routing=com.azure.cosmos.kafka.connect

</javaModulesSurefireArgLine>
<doclintMissingInclusion>-</doclintMissingInclusion>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -447,6 +447,7 @@ private Mono<Map<FeedRange, KafkaCosmosChangeFeedState>> getEffectiveContinuatio
containerFeedRange,
this.getContinuationStateFromOffset(
feedRangeContinuationTopicOffset,
containerFeedRange,
containerFeedRange));

return Mono.just(effectiveContinuationMap);
Expand Down Expand Up @@ -481,11 +482,14 @@ private Mono<Map<FeedRange, KafkaCosmosChangeFeedState>> getEffectiveContinuatio
);

if (continuationTopicOffset == null) {
effectiveContinuationMap.put(overlappedFeedRangesFromOffset.get(0), null);
effectiveContinuationMap.put(containerFeedRange, null);
} else {
effectiveContinuationMap.put(
containerFeedRange,
this.getContinuationStateFromOffset(continuationTopicOffset, containerFeedRange));
this.getContinuationStateFromOffset(
continuationTopicOffset,
overlappedFeedRangesFromOffset.get(0),
containerFeedRange));
}

return Mono.just(effectiveContinuationMap);
Expand All @@ -507,6 +511,7 @@ private Mono<Map<FeedRange, KafkaCosmosChangeFeedState>> getEffectiveContinuatio
overlappedRangeFromOffset,
this.getContinuationStateFromOffset(
this.kafkaOffsetStorageReader.getFeedRangeContinuationOffset(databaseName, containerRid, overlappedRangeFromOffset),
overlappedRangeFromOffset,
overlappedRangeFromOffset));
}
}
Expand All @@ -522,15 +527,16 @@ private Mono<Map<FeedRange, KafkaCosmosChangeFeedState>> getEffectiveContinuatio

private KafkaCosmosChangeFeedState getContinuationStateFromOffset(
FeedRangeContinuationTopicOffset feedRangeContinuationTopicOffset,
FeedRange feedRange) {
FeedRange offsetFeedRange,
FeedRange targetFeedRange) {

KafkaCosmosChangeFeedState changeFeedState =
KafkaCosmosChangeFeedState offsetState =
new KafkaCosmosChangeFeedState(
feedRangeContinuationTopicOffset.getResponseContinuation(),
feedRange,
offsetFeedRange,
feedRangeContinuationTopicOffset.getItemLsn());

return changeFeedState;
return offsetState.extractForFeedRange(targetFeedRange);
}

private List<FeedRange> getFeedRanges(CosmosContainerProperties containerProperties) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -407,7 +407,11 @@ private Mono<Boolean> handleFeedRangeGone(FeedRangeTaskUnit feedRangeTaskUnit) {
.getOverlappingFeedRanges(container, feedRangeTaskUnit.getFeedRange(), true)
.flatMap(overlappedRanges -> {

if (overlappedRanges.size() == 1) {
if (overlappedRanges.isEmpty()) {
return Mono.error(
new IllegalStateException(
"No overlapping feed ranges found for " + feedRangeTaskUnit.getFeedRange()));
Comment thread
tvaron3 marked this conversation as resolved.
} else if (overlappedRanges.size() == 1) {
// merge happens
LOGGER.info(
"FeedRange {} is merged into {}, but we will continue polling data from feedRange {}",
Expand Down Expand Up @@ -478,7 +482,7 @@ private KafkaCosmosChangeFeedState getChildRangeChangeFeedState(
KafkaCosmosChangeFeedState parent,
FeedRange feedRange) {
return parent == null
? null : new KafkaCosmosChangeFeedState(parent.getResponseContinuation(), feedRange);
? null : parent.extractForFeedRange(feedRange);
}

private CosmosChangeFeedRequestOptions getChangeFeedRequestOptions(FeedRangeTaskUnit feedRangeTaskUnit) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

package com.azure.cosmos.kafka.connect.implementation.source;

import com.azure.cosmos.implementation.ImplementationBridgeHelpers;
import com.azure.cosmos.implementation.apachecommons.lang.StringUtils;
import com.azure.cosmos.models.FeedRange;
import com.fasterxml.jackson.core.JsonGenerator;
Expand Down Expand Up @@ -52,6 +53,18 @@ public String getItemLsn() {
return itemLsn;
}

public KafkaCosmosChangeFeedState extractForFeedRange(FeedRange feedRange) {
checkNotNull(feedRange, "Argument 'feedRange' can not be null");

String projectedContinuation = ImplementationBridgeHelpers
.CosmosChangeFeedRequestOptionsHelper
.getCosmosChangeFeedRequestOptionsAccessor()
.extractContinuationForFeedRange(this.responseContinuation, feedRange);
String projectedItemLsn = this.targetRange.equals(feedRange) ? this.itemLsn : null;

return new KafkaCosmosChangeFeedState(projectedContinuation, feedRange, projectedItemLsn);
}

@Override
public boolean equals(Object o) {
if (this == o) {
Expand Down Expand Up @@ -108,7 +121,11 @@ public KafkaCosmosChangeFeedState deserialize(
final JsonNode rootNode = jsonParser.getCodec().readTree(jsonParser);
String continuationState = rootNode.get("responseContinuation").asText();
FeedRange targetRange = FeedRange.fromString(rootNode.get("targetRange").asText());
String continuationLsn = rootNode.get("itemLsn").asText();
String continuationLsn = null;
if (rootNode.hasNonNull("itemLsn")) {
String itemLsnValue = rootNode.get("itemLsn").asText();
continuationLsn = "null".equals(itemLsnValue) ? null : itemLsnValue;
}
return new KafkaCosmosChangeFeedState(continuationState, targetRange, continuationLsn);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,19 @@ public void getTaskConfigsAfterSplit() throws JsonProcessingException {
multiPartitionContainer.getResourceId(),
FeedRange.forFullRange());

List<FeedRange> currentFeedRanges = cosmosAsyncClient
.getDatabase(databaseName)
.getContainer(multiPartitionContainer.getId())
.getFeedRanges()
.block();
List<CompositeContinuationToken> childContinuationTokens = new ArrayList<>();
for (int i = 0; i < currentFeedRanges.size(); i++) {
childContinuationTokens.add(
new CompositeContinuationToken(
String.valueOf(100 + i),
((FeedRangeEpkImpl) currentFeedRanges.get(i)).getRange()));
}

String initialContinuationState = new ChangeFeedStateV1(
multiPartitionContainer.getResourceId(),
FeedRangeEpkImpl.forFullRange(),
Expand All @@ -380,10 +393,10 @@ public void getTaskConfigsAfterSplit() throws JsonProcessingException {
FeedRangeContinuation.create(
multiPartitionContainer.getResourceId(),
FeedRangeEpkImpl.forFullRange(),
Arrays.asList(new CompositeContinuationToken("1", FeedRangeEpkImpl.forFullRange().getRange())))).toString();
childContinuationTokens)).toString();

FeedRangeContinuationTopicOffset feedRangeContinuationTopicOffset =
new FeedRangeContinuationTopicOffset(initialContinuationState, "1"); // using the same itemLsn as in the continuationToken
new FeedRangeContinuationTopicOffset(initialContinuationState, "100");
Map<Map<String, Object>, Map<String, Object>> initialOffsetMap = new HashMap<>();
initialOffsetMap.put(
FeedRangeContinuationTopicPartition.toMap(feedRangeContinuationTopicPartition),
Expand Down Expand Up @@ -844,11 +857,13 @@ private List<FeedRangeTaskUnit> getFeedRangeTaskUnits(
KafkaCosmosChangeFeedState kafkaCosmosChangeFeedState = null;
if (StringUtils.isNotEmpty(continuationState)) {
ChangeFeedState changeFeedState = ChangeFeedStateV1.fromString(continuationState);
ChangeFeedState projectedState = changeFeedState.extractForEffectiveRange(
((FeedRangeEpkImpl) feedRange).getRange());
kafkaCosmosChangeFeedState =
new KafkaCosmosChangeFeedState(
continuationState,
projectedState.toString(),
feedRange,
changeFeedState.getContinuation().getCurrentContinuationToken().getToken());
null);
}

return new FeedRangeTaskUnit(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,29 @@

import com.azure.cosmos.CosmosAsyncClient;
import com.azure.cosmos.CosmosAsyncContainer;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedMode;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStartFromInternal;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedState;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStateV1;
import com.azure.cosmos.implementation.feedranges.FeedRangeContinuation;
import com.azure.cosmos.implementation.feedranges.FeedRangeEpkImpl;
import com.azure.cosmos.implementation.query.CompositeContinuationToken;
import com.azure.cosmos.implementation.routing.PartitionKeyInternalHelper;
import com.azure.cosmos.kafka.connect.InMemoryStorageReader;
import com.azure.cosmos.kafka.connect.KafkaCosmosTestConfigurations;
import com.azure.cosmos.kafka.connect.KafkaCosmosTestSuiteBase;
import com.azure.cosmos.kafka.connect.TestItem;
import com.azure.cosmos.kafka.connect.implementation.CosmosClientCache;
import com.azure.cosmos.kafka.connect.implementation.CosmosClientCacheItem;
import com.azure.cosmos.models.CosmosChangeFeedRequestOptions;
import com.azure.cosmos.models.CosmosContainerProperties;
import com.azure.cosmos.models.CosmosItemRequestOptions;
import com.azure.cosmos.models.CosmosQueryRequestOptions;
import com.azure.cosmos.models.FeedRange;
import com.azure.cosmos.models.FeedResponse;
import com.azure.cosmos.models.ModelBridgeInternal;
import com.azure.cosmos.models.PartitionKey;
import com.azure.cosmos.models.PartitionKeyDefinition;
import com.azure.cosmos.models.ThroughputProperties;
import com.azure.cosmos.models.ThroughputResponse;
import org.apache.kafka.connect.data.Struct;
Expand All @@ -26,6 +39,7 @@

import java.util.ArrayList;
import java.util.Arrays;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -226,6 +240,99 @@ public void poll_splitWhenStartFeedRangeTask() {
}
}
}

@Test(groups = { "kafka" }, timeOut = 60 * TIMEOUT)
public void pollWithDivergentChildContinuations() {
Map<String, String> sourceConfigMap = new HashMap<>();
sourceConfigMap.put("azure.cosmos.account.endpoint", KafkaCosmosTestConfigurations.HOST);
sourceConfigMap.put("azure.cosmos.account.key", KafkaCosmosTestConfigurations.MASTER_KEY);
sourceConfigMap.put("azure.cosmos.source.database.name", databaseName);
sourceConfigMap.put(
"azure.cosmos.source.containers.includedList",
Arrays.asList(multiPartitionContainerName).toString());
sourceConfigMap.put("azure.cosmos.source.task.id", UUID.randomUUID().toString());

CosmosSourceConfig sourceConfig = new CosmosSourceConfig(sourceConfigMap);
CosmosClientCacheItem clientItem =
CosmosClientCache.getCosmosClient(sourceConfig.getAccountConfig(), "pollWithDivergentChildContinuations");
CosmosSourceTask sourceTask = null;

try {
CosmosAsyncContainer container = clientItem
.getClient()
.getDatabase(databaseName)
.getContainer(multiPartitionContainerName);
CosmosContainerProperties containerProperties = container.read().block().getProperties();
List<FeedRangeEpkImpl> childRanges = container.getFeedRanges().block().stream()
.map(range -> (FeedRangeEpkImpl) range)
.sorted(Comparator.comparing(range -> range.getRange().getMin()))
.collect(Collectors.toList());
assertThat(childRanges.size()).isGreaterThanOrEqualTo(2);

FeedRangeEpkImpl childA = childRanges.get(0);
FeedRangeEpkImpl childB = childRanges.get(1);
ChangeFeedState baselineB = snapshotNow(container, childB);

ChangeFeedState advancedA = snapshotNow(container, childA);
int writesToChildA = 0;
while (numericToken(advancedA) <= numericToken(baselineB) && writesToChildA < 2_000) {
for (int i = 0; i < 50; i++) {
container.createItem(new TestItem(
UUID.randomUUID().toString(),
partitionKeyForRange("advance-a", childA, containerProperties.getPartitionKeyDefinition()),
"advance-a")).block();
writesToChildA++;
}
advancedA = snapshotNow(container, childA);
}
assertThat(numericToken(advancedA)).isGreaterThan(numericToken(baselineB));

ChangeFeedState parentState = new ChangeFeedStateV1(
containerProperties.getResourceId(),
FeedRangeEpkImpl.forFullRange(),
ChangeFeedMode.INCREMENTAL,
ChangeFeedStartFromInternal.createFromNow(),
FeedRangeContinuation.create(
containerProperties.getResourceId(),
FeedRangeEpkImpl.forFullRange(),
Arrays.asList(
new CompositeContinuationToken(rawToken(advancedA), childA.getRange()),
new CompositeContinuationToken(rawToken(baselineB), childB.getRange()))));

String markerId = UUID.randomUUID().toString();
container.createItem(
new TestItem(
markerId,
partitionKeyForRange("marker-b", childB, containerProperties.getPartitionKeyDefinition()),
"marker-b")).block();

FeedRangeTaskUnit feedRangeTaskUnit = new FeedRangeTaskUnit(
databaseName,
multiPartitionContainerName,
containerProperties.getResourceId(),
childB,
new KafkaCosmosChangeFeedState(parentState.toString(), childB),
multiPartitionContainerName);
Map<String, String> taskConfigMap = sourceConfig.originalsStrings();
taskConfigMap.putAll(
CosmosSourceTaskConfig.getFeedRangeTaskUnitsConfigMap(Arrays.asList(feedRangeTaskUnit)));

sourceTask = new CosmosSourceTask();
sourceTask.initialize(new TestSourceTaskContext(taskConfigMap));
sourceTask.start(taskConfigMap);

List<SourceRecord> records = sourceTask.poll();
assertThat(records.stream().anyMatch(record ->
markerId.equals(((Struct) record.value()).get("id").toString()))).isTrue();
} finally {
if (sourceTask != null) {
sourceTask.stop();
}
CosmosClientCache.releaseCosmosClient(clientItem.getClientConfig());
clientItem.getClient().close();
}
}

@Test(groups = { "kafka", "kafka-emulator" }, timeOut = TIMEOUT)
public void pollWithSpecificFeedRange() {
// Test only items belong to the feedRange defined in the feedRangeTaskUnit will be returned
Expand Down Expand Up @@ -513,6 +620,43 @@ private List<TestItem> createItems(
return testItems;
}

private ChangeFeedState snapshotNow(CosmosAsyncContainer container, FeedRangeEpkImpl feedRange) {
FeedResponse<com.fasterxml.jackson.databind.JsonNode> response = container
.queryChangeFeed(
CosmosChangeFeedRequestOptions.createForProcessingFromNow(feedRange),
com.fasterxml.jackson.databind.JsonNode.class)
.byPage()
.next()
.block();
return ChangeFeedState.fromString(response.getContinuationToken());
}

private String partitionKeyForRange(
String prefix,
FeedRangeEpkImpl targetRange,
PartitionKeyDefinition partitionKeyDefinition) {

for (int i = 0; i < 100_000; i++) {
String id = prefix + "-" + UUID.randomUUID();
String effectivePartitionKey = PartitionKeyInternalHelper.getEffectivePartitionKeyString(
ModelBridgeInternal.getPartitionKeyInternal(new PartitionKey(id)),
partitionKeyDefinition);
if (targetRange.getRange().contains(effectivePartitionKey)) {
return id;
}
}

throw new IllegalStateException("Unable to generate an id for feed range " + targetRange);
}

private long numericToken(ChangeFeedState state) {
return Long.parseLong(rawToken(state).replace("\"", ""));
}

private String rawToken(ChangeFeedState state) {
return state.getContinuation().getCurrentContinuationToken().getToken();
}

public static class TestSourceTaskContext implements SourceTaskContext {
private final Map<String, String> map;
private final OffsetStorageReader offsetStorageReader;
Expand Down
Loading
Loading