You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Scala实现BigQuery数据查询直传S3的技术求助

Directly Export BigQuery Query Results to S3 (Bypassing GCS & gsutil)

Got it, let’s get this sorted for you. As a BigQuery newbie, skipping GCS and gsutil to move query results straight to S3 is totally feasible with Scala—here’s a step-by-step solution including a working code example.


Core Approach

Instead of routing data through GCS, we’ll:

  • Use the official BigQuery Scala/Java client to run your SQL query and fetch results.
  • Stream the query output directly to an S3 object using the AWS SDK for Java (fully compatible with Scala).
  • Avoid loading all data into memory by using streaming writes (critical for large datasets).

Required Dependencies

Add these to your build.sbt file to pull in the necessary libraries:

libraryDependencies ++= Seq(
  "com.google.cloud" % "google-cloud-bigquery" % "2.30.0",
  "software.amazon.awssdk" % "s3" % "2.25.0",
  "com.opencsv" % "opencsv" % "5.6" // For CSV formatting (swap if using Parquet/JSON)
)

Working Scala Code Example

This example runs a BigQuery query, formats results as CSV, and streams them directly to S3:

import com.google.cloud.bigquery._
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider
import software.amazon.awssdk.regions.Region
import software.amazon.awssdk.services.s3.S3Client
import software.amazon.awssdk.services.s3.model.PutObjectRequest
import com.opencsv.CSVWriter
import java.io.{ByteArrayInputStream, ByteArrayOutputStream}
import scala.collection.JavaConverters._

object BigQueryToS3 {
  def main(args: Array[String]): Unit = {
    // 1. Initialize BigQuery client
    val bigQuery = BigQueryOptions.getDefaultInstance.getService

    // 2. Define your target query
    val query = """
      SELECT column1, column2, column3 
      FROM `your-gcp-project.your-dataset.your-table`
      WHERE date >= '2024-01-01'
    """.trim

    // 3. Execute query and wait for completion
    val queryJob = bigQuery.query(QueryJobConfiguration.newBuilder(query).build())
    while (!queryJob.isDone) Thread.sleep(1000)
    val resultSet = queryJob.getQueryResults

    // 4. Stream results to CSV (avoids loading all data into memory)
    val outputStream = new ByteArrayOutputStream()
    val csvWriter = new CSVWriter(new java.io.OutputStreamWriter(outputStream))
    
    // Write header row
    val columnNames = resultSet.getSchema.getFields.asScala.map(_.getName).toArray
    csvWriter.writeNext(columnNames)
    
    // Write data rows
    resultSet.iterateAll().asScala.foreach(row => {
      val rowValues = row.getValues.asScala.map(_.getValue.toString).toArray
      csvWriter.writeNext(rowValues)
    })
    
    csvWriter.close()
    val csvData = outputStream.toByteArray
    outputStream.close()

    // 5. Initialize S3 client
    val s3Client = S3Client.builder()
      .region(Region.US_EAST_1) // Update to your bucket's region
      .credentialsProvider(DefaultCredentialsProvider.create())
      .build()

    // 6. Upload CSV to S3
    val bucketName = "your-s3-bucket-name"
    val s3FileKey = "bigquery-results/2024-q1-data.csv" // Path inside your bucket

    val uploadRequest = PutObjectRequest.builder()
      .bucket(bucketName)
      .key(s3FileKey)
      .contentType("text/csv")
      .build()

    s3Client.putObject(uploadRequest, software.amazon.awssdk.core.sync.RequestBody.fromBytes(csvData))

    // Cleanup and confirm success
    s3Client.close()
    println(s"✅ Results uploaded successfully to s3://$bucketName/$s3FileKey")
  }
}

Key Notes & Troubleshooting

  • Authentication:
    • BigQuery: Run gcloud auth application-default login locally, or set the GOOGLE_APPLICATION_CREDENTIALS environment variable to your service account key file path.
    • AWS: Ensure your ~/.aws/credentials file has valid keys, or set AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY environment variables.
  • Large Datasets: For millions of rows, replace the in-memory CSV stream with a temporary file or use S3’s multipart upload API to avoid memory overload.
  • Data Formats: To use Parquet instead of CSV, integrate libraries like org.apache.parquet to format the stream before uploading.
  • Permissions: Your BigQuery service account needs bigquery.jobs.create and bigquery.tables.getData permissions. Your AWS IAM role/user needs s3:PutObject access to the target bucket.

内容的提问来源于stack exchange,提问作者Tech Guy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 06:37:41