Skip to content

Commit b73c1fa

Browse files
committed
spark: use UC Delta Rest Catalog API credentials for external createTable
1 parent fbe7989 commit b73c1fa

2 files changed

Lines changed: 82 additions & 0 deletions

File tree

spark/src/main/scala/org/apache/spark/sql/delta/catalog/DeltaCatalogClient.scala

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -157,6 +157,21 @@ private class DeltaCatalogClient private (
157157
schemaName,
158158
tableName,
159159
staging.getTableId.toString)))
160+
case (CatalogTableType.EXTERNAL, Some(externalLocation))
161+
if isCloudScheme(externalLocation.getScheme) =>
162+
val locationText = externalLocation.toString
163+
// External create must write the initial _delta_log, so READ fallback would be wrong.
164+
val credentials =
165+
client.getTemporaryPathCredentials(locationText, CredentialOperation.READ_WRITE)
166+
Some(PreparedUCDeltaRestCatalogApiCreate(
167+
location = externalLocation,
168+
tableProperties = Map.empty,
169+
storageProperties = toPathCredentialProperties(
170+
locationText,
171+
getStorageCredentials(credentials),
172+
externalLocation.getScheme,
173+
PathOperation.PATH_CREATE_TABLE,
174+
credentialContext)))
160175
case _ =>
161176
None
162177
}
@@ -566,6 +581,33 @@ private[delta] object DeltaCatalogClient {
566581
}
567582
}
568583

584+
private def toPathCredentialProperties(
585+
location: String,
586+
storageCredentials: Seq[StorageCredential],
587+
locationScheme: String,
588+
pathOperation: PathOperation,
589+
credentialContext: Option[UCDeltaRestCatalogApiCredentialContext]): Map[String, String] = {
590+
cloudCredentialProperties(
591+
location,
592+
storageCredentials,
593+
locationScheme,
594+
credentialContext) { (context, credential) =>
595+
CredPropsUtil.createPathCredProps(
596+
context.renewCredentialEnabled,
597+
context.credScopedFsEnabled,
598+
context.hadoopConf,
599+
locationScheme.toLowerCase(Locale.ROOT),
600+
context.uri,
601+
context.tokenProvider,
602+
location,
603+
pathOperation,
604+
DeltaStorageCredentialUtil.toTemporaryCredentials(
605+
toUnityCatalogStorageCredential(credential)))
606+
.asScala
607+
.toMap
608+
}
609+
}
610+
569611
private def toTableProperties(staging: StagingTableResponse): Map[String, String] = {
570612
val stagingTableId = staging.getTableId.toString
571613
val requiredProperties = Option(staging.getRequiredProperties)

spark/src/test/scala/org/apache/spark/sql/delta/catalog/DeltaCatalogClientSuite.scala

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -785,6 +785,46 @@ class DeltaCatalogClientSuite
785785
}
786786
}
787787

788+
test("prepareCreateTable uses temporary path credentials for external cloud tables") {
789+
val location = "s3://bucket/external/tbl"
790+
configHandler = exchange => {
791+
assert(queryParams(exchange)("catalog") === "uc")
792+
sendJson(exchange, 200,
793+
"""{
794+
| "endpoints": [
795+
| "GET /v1/catalogs/{catalog}/schemas/{schema}/tables/{table}",
796+
| "GET /v1/catalogs/{catalog}/schemas/{schema}/tables/{table}/credentials",
797+
| "GET /v1/temporary-path-credentials"
798+
| ],
799+
| "protocol-version": "1.0"
800+
|}""".stripMargin)
801+
}
802+
pathCredentialsHandler = exchange => {
803+
credentialRequestCount += 1
804+
assert(queryParams(exchange) === Map(
805+
"location" -> location,
806+
"operation" -> "READ_WRITE"))
807+
sendJson(exchange, 200, s3CredentialsResponseJson(location, "READ_WRITE"))
808+
}
809+
810+
val prepared = withUCDeltaRestCatalogApi { catalog =>
811+
catalog.prepareCreateTable(
812+
Identifier.of(Array("default"), "tbl"),
813+
CatalogTableType.EXTERNAL,
814+
location = Some(java.net.URI.create(location))).get
815+
}
816+
817+
assert(credentialRequestCount === 1)
818+
assert(prepared.location.toString === location)
819+
assert(prepared.tableProperties.isEmpty)
820+
assert(prepared.storageProperties(S3ACredentialsProviderKey) ===
821+
AwsVendedTokenProviderClass)
822+
assert(prepared.storageProperties(UCCredentialsTypeKey) === UCCredentialsTypePathValue)
823+
assert(prepared.storageProperties(UCPathOperationKey) ===
824+
"PATH_CREATE_TABLE")
825+
assert(prepared.storageProperties(UCPathKey) === location)
826+
}
827+
788828
private def loadWithUCDeltaRestCatalogApi(): V1Table = {
789829
withUCDeltaRestCatalogApi { catalog =>
790830
catalog.loadTable(Identifier.of(Array("default"), "tbl")).get.asInstanceOf[V1Table]

0 commit comments

Comments
 (0)