Azure / Azure/azure-sdk-for-java

TransientIOErrorsRetryingIterator trigger one more roundtrip when using Cosmos Spark connector to query CosmosDB data

Open
#47,777 7 comments 0 reactions 2 assignees Claimed by @xinlian12 View on GitHub
Client Cosmos customer-reported needs-team-attention question Service Attention
Dominant language
Java
Stars
2.6k
Forks
2.2k
Avg merge
2d 8h
Merged PRs (30d)
178

Description

When connecting to CosmosDB to query data through Spark connector, e.g. the total item count = 20k, setting maxItemCount = 5k in the client code, the round trip is supposed to be 4, however, it often somehow triggers 5 round trips, which brings extra RU consumption.

We found that by commenting out 'cancelOn' and recompiled the jar file for testing, the issue never reoccurred, roundtrip count is always 4, which is correct count. Is this a BUG? Or any particular purpose of the below code?
override def close(): Unit = {
lastPagedFlux.getAndSet(None) match {
case Some(oldPagedFlux) => oldPagedFlux.cancelOn(Schedulers.boundedElastic()).onErrorComplete().subscribe().dispose()
case None =>
}
}

https://github.com/Azure/azure-sdk-for-java/blob/429ad868ade2458fd283ea24161c13d778d3d97d/sdk/cosmos/azure-cosmos-spark_3/src/main/scala/com/azure/cosmos/spark/TransientIOErrorsRetryingIterator.scala#L283

Below is the code I use in data bricks:
-----------------------------------------
from pyspark.sql.types import StructType, StructField, StringType

cosmos_schema = StructType([
StructField('createddatetime', StringType(), True),
StructField('closeddatetime', StringType(), True),
StructField('id', StringType(), False),
StructField('owner', StringType(), True),
StructField('ownerregion', StringType(), True),
StructField('ownermanager', StringType(), True),
StructField('casenumber', StringType(), True)
])

cosmos_config = {
"spark.cosmos.accountEndpoint": COSMOS_ENDPOINT,
"spark.cosmos.accountkey": CosmosMasterKey,
"spark.cosmos.database": COSMOS_DATABASE,
"spark.cosmos.container": COSMOS_CONTAINER,
"spark.cosmos.diagnostics": "simple",
'spark.cosmos.read.inferSchema.enabled': 'false',
"spark.cosmos.read.customQuery": customQuery,
"spark.cosmos.read.maxItemCount": 5000
}

sc.setLogLevel("DEBUG")

df = spark.read.format("cosmos.oltp").options(**cosmos_config).schema(cosmos_schema).load()
total_item_count=df.count()
print(f'Total documents retrieved: {total_item_count}')

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.