feat: add auth to datalake (#7852)

Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com>
This commit is contained in:
Alexander Onnikov
2025-01-31 20:51:20 +07:00
committed by GitHub
parent b11e17aee0
commit e11e0eb0dc
7 changed files with 117 additions and 45 deletions
+2 -1
View File
@@ -1279,10 +1279,11 @@ export function devTool (
blobsSize: workspace.backupInfo?.blobsSize ?? 0
})
const workspaceId = getWorkspaceId(workspace.workspace)
const token = generateToken(systemAccountEmail, workspaceId)
for (const config of storages) {
const storage = new S3Service(config)
await copyToDatalake(toolCtx, workspaceId, config, storage, datalake, params)
await copyToDatalake(toolCtx, workspaceId, config, storage, datalake, token, params)
}
}
})
+9 -6
View File
@@ -271,6 +271,7 @@ export async function copyToDatalake (
config: S3Config,
adapter: S3Service,
datalake: DatalakeClient,
token: string,
params: CopyDatalakeParams
): Promise<void> {
console.log('copying from', config.name, 'concurrency:', params.concurrency)
@@ -311,7 +312,7 @@ export async function copyToDatalake (
let cursor: string | undefined = ''
let hasMore = true
while (hasMore) {
const res = await datalake.listObjects(ctx, workspaceId, cursor, 1000)
const res = await datalake.listObjects(ctx, token, workspaceId, cursor, 1000)
cursor = res.cursor
hasMore = res.cursor !== undefined
for (const blob of res.blobs) {
@@ -349,7 +350,7 @@ export async function copyToDatalake (
ctx,
5,
async () => {
await copyBlobToDatalake(ctx, workspaceId, blob, config, adapter, datalake)
await copyBlobToDatalake(ctx, workspaceId, blob, config, adapter, datalake, token)
processedCnt += 1
processedSize += blob.size
},
@@ -378,7 +379,8 @@ export async function copyBlobToDatalake (
blob: ListBlobResult,
config: S3Config,
adapter: S3Service,
datalake: DatalakeClient
datalake: DatalakeClient,
token: string
): Promise<void> {
const objectName = blob._id
if (blob.size < 1024 * 1024 * 64) {
@@ -390,7 +392,7 @@ export async function copyBlobToDatalake (
const url = concatLink(endpoint, `${bucketId}/${objectId}`)
const params = { url, accessKeyId, secretAccessKey, region }
await datalake.uploadFromS3(ctx, workspaceId, objectName, params)
await datalake.uploadFromS3(ctx, token, workspaceId, objectName, params)
} else {
// Handle huge file
const stat = await adapter.stat(ctx, workspaceId, objectName)
@@ -405,7 +407,7 @@ export async function copyBlobToDatalake (
const readable = await adapter.get(ctx, workspaceId, objectName)
try {
console.log('uploading huge blob', objectName, Math.round(stat.size / 1024 / 1024), 'MB')
await uploadMultipart(ctx, datalake, workspaceId, objectName, readable, metadata)
await uploadMultipart(ctx, token, datalake, workspaceId, objectName, readable, metadata)
console.log('done', objectName)
} finally {
readable.destroy()
@@ -416,6 +418,7 @@ export async function copyBlobToDatalake (
function uploadMultipart (
ctx: MeasureContext,
token: string,
datalake: DatalakeClient,
workspaceId: WorkspaceId,
objectName: string,
@@ -446,7 +449,7 @@ function uploadMultipart (
stream.pipe(passthrough)
datalake
.uploadMultipart(ctx, workspaceId, objectName, passthrough, metadata)
.uploadMultipart(ctx, token, workspaceId, objectName, passthrough, metadata)
.then(() => {
cleanup()
resolve()