【问题标题】:Create Avro file directly to Google Cloud Storage将 Avro 文件直接创建到 Google Cloud Storage
【发布时间】:2019-12-02 19:30:19
【问题描述】:

我想跳过在本地创建 avro 文件并将其直接上传到 Google Cloud Storage 的步骤。

我检查了 blob.upload from_string 选项,但老实说,我不知道应该替换什么以应用于我的代码。而且我不知道这是否是满足我需要的最佳方式。有了它,我可以通过将脚本包含在 docker 映像中来构建更现代的管道。

这可以根据下面的脚本以某种方式完成:

import csv
import base64
import json
import io
import avro.schema
import avro.io
from avro.datafile import DataFileReader, DataFileWriter
import math
import os
import gcloud
from gcloud import storage
from google.cloud import bigquery
from oauth2client.client import GoogleCredentials
from datetime import datetime, timedelta
import numpy as np

try:
    script_path = os.path.dirname(os.path.abspath(__file__)) + "/"
except:
    script_path = "C:\\Users\\me\\Documents\\Keys\\key.json"

#Bigquery Credentials and settings
os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = script_path

folder = str((datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d'))
bucket_name = 'gs://new_bucket/table/*.csv'
dataset = 'dataset'
tabela = 'table'

schema = avro.schema.Parse(open("C:\\Users\\me\\schema_table.avsc", "rb").read())  

writer = DataFileWriter(open("C:\\Users\\me\\table_register.avro", "wb"), avro.io.DatumWriter(), schema)


def insert_bigquery(target_uri, dataset_id, table_id):
    bigquery_client = bigquery.Client()
    dataset_ref = bigquery_client.dataset(dataset_id)
    job_config = bigquery.LoadJobConfig()
    job_config.schema = [
        bigquery.SchemaField('id','STRING',mode='REQUIRED')
    ]
    job_config.source_format = bigquery.SourceFormat.CSV
    job_config.field_delimiter = ";"
    uri = target_uri
    load_job = bigquery_client.load_table_from_uri(
        uri,
        dataset_ref.table(table_id),
        job_config=job_config
        )
    print('Starting job {}'.format(load_job.job_id))
    load_job.result()
    print('Job finished.')

#insert_bigquery(bucket_name, dataset, tabela)

def get_data_from_bigquery():
    """query bigquery to get data to import to PSQL"""
    bq = bigquery.Client()
    #Busca IDs
    query = """SELECT id FROM dataset.base64_data"""
    query_job = bq.query(query)
    data = query_job.result()
    rows = list(data)
    return rows

a = get_data_from_bigquery()
length = len(a) 
line_count = 0

for row in range(length):
    bytes = base64.b64decode(str(a[row][0]))
    bytes = bytes[5:]
    buf = io.BytesIO(bytes)
    decoder = avro.io.BinaryDecoder(buf)
    rec_reader = avro.io.DatumReader(avro.schema.Parse(open("C:\\Users\\me\\schema_table.avsc").read()))
    out=rec_reader.read(decoder)
    writer.append(out)
writer.close()

def upload_blob(bucket_name, source_file_name, destination_blob_name):
    storage_client = storage.Client()
    bucket = storage_client.get_bucket(bucket_name)
    blob = bucket.blob("insert_transfer/" + destination_blob_name)
    blob.upload_from_filename(source_file_name)
    print('File {} uploaded to {}'.format(
        source_file_name,
        destination_blob_name
    ))

upload_blob('new_bucket', 'C:\\Users\\me\\table_register.avro', 'table_register.avro')

【问题讨论】:

  • 文件大小是多少?如果它们不是很大,这意味着小于可用的可用内存,则在内存中构建 avro 文件并从字符串上传 blob。否则,您将不得不将数据流式传输到 Cloud Storage,这并不容易,这意味着无法使用 Google 库。
  • @JohnHanley 你能给我看一个基于这个脚本的例子吗?我的文件不大!
  • 抱歉,我不拥有我为 Avro 编写的代码。我的评论是给你一个可能的解决方案的想法。

标签: python google-cloud-storage avro


【解决方案1】:

我看到了您的脚本,并且可以看到您正在从 BigQuery 获取数据。我可以向您确认,我重现了您的场景,并且可以直接将数据从 BigQuery 导出到 Google Cloud Storage,而无需在本地创建 avro 文件。

我建议您查看here,其中描述了如何将表数据从 BigQuery 导出到 Google Cloud Storage。以下是要遵循的步骤:

  1. 在您的 Cloud Console 中打开 BigQuery 网页界面。
  2. 在导航面板的“资源”部分,展开您的项目并单击 您的数据集以扩展它。查找并单击包含您的数据的表格 正在导出。
  3. 在窗口右侧,点击导出,然后选择导出到云存储
  4. 在“导出到云存储”对话框中:
    • 对于选择云存储位置,浏览存储桶。
    • 对于导出格式,选择导出数据的格式,在您的特定 情况下,选择“Avro”。
    • 点击导出。

不过,也有可能使用 Python 来实现。我建议你看看here

我希望这种方法对您有用。

【讨论】:

  • 我需要在转换为 avro 之前进行 BigQuery 数据转换。在那种情况下,解决方案不适合我!
  • 由于您在从 BigQuery 检索文件后尝试进行某种转换,我最好的解决方案是在转换为 Avro 后下载文件,然后上传到您的 Google Cloud Storage Bucket .最后,一旦上传到您的 Cloud Storage Bucket,您就可以将其删除。
猜你喜欢
  • 1970-01-01
  • 2021-01-22
  • 1970-01-01
  • 1970-01-01
  • 2018-11-18
  • 2020-01-15
  • 2018-06-25
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多