【问题标题】:Error trying to write CSV file to Google Cloud Storage from Dataflow pipeline尝试从 Dataflow 管道将 CSV 文件写入 Google Cloud Storage 时出错
【发布时间】:2022-01-12 17:16:41
【问题描述】:

我正在构建一个 Dataflow 管道,该管道从我的 Cloud Storage 存储桶中读取一个 CSV 文件(包含 250,000 行),修改每行的值,然后将修改后的内容写入同一存储桶中的新 CSV。使用下面的代码,我可以读取和修改原始文件的内容,但是当我尝试在 GCS 中写入新文件的内容时,我收到以下错误:

google.api_core.exceptions.TooManyRequests: 429 POST https://storage.googleapis.com/upload/storage/v1/b/my-bucket/o?uploadType=multipart: {
  "error": {
    "code": 429,
    "message": "The rate of change requests to the object my-bucket/product-codes/URL_test_codes.csv exceeds the rate limit. Please reduce the rate of create, update, and delete requests.",
    "errors": [
      {
        "message": "The rate of change requests to the object my-bucket/product-codes/URL_test_codes.csv exceeds the rate limit. Please reduce the rate of create, update, and delete requests.",
        "domain": "usageLimits",
        "reason": "rateLimitExceeded"
      }
    ]
  }
}
: ('Request failed with status code', 429, 'Expected one of', <HTTPStatus.OK: 200>) [while running 'Store Output File']

我在 Dataflow 中的代码:

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
import traceback
import sys
import pandas as pd
from cryptography.fernet import Fernet
import google.auth
from google.cloud import storage

fernet_secret = 'aD4t9MlsHLdHyuFKhoyhy9_eLKDfe8eyVSD3tu8KzoP='
bucket = 'my-bucket'
inputFile = f'gs://{bucket}/product-codes/test_codes.csv'
outputFile = 'product-codes/URL_test_codes.csv'

#Pipeline Logic
def product_codes_pipeline(project, env, region='us-central1'):
    options = PipelineOptions(
        streaming=False,
        project=project,
        region=region,
        staging_location="gs://my-bucket-dataflows/Templates/staging",
        temp_location="gs://my-bucket-dataflows/Templates/temp",
        template_location="gs://my-bucket-dataflows/Templates/Generate_Product_Codes.py",
        subnetwork='https://www.googleapis.com/compute/v1/projects/{}/regions/us-central1/subnetworks/{}-private'.format(project, env)
    )
    
    # Transform function
    def genURLs(code):
        f = Fernet(fernet_secret)
        encoded = code.encode()
        encrypted = f.encrypt(encoded)
        decrypted = f.decrypt(encrypted.decode().encode())
        decoded = decrypted.decode()
        if code != decoded:
            print(f'Error: Code {code} and decoded code {decoded} do not match')
            sys.exit(1)
        url = 'https://some-url.com/redeem/product-code=' + encrypted.decode()
        return url
    
    class WriteCSVFIle(beam.DoFn):
        def __init__(self, bucket_name):
            self.bucket_name = bucket_name

        def start_bundle(self):
            self.client = storage.Client()

        def process(self, urls):
            df = pd.DataFrame([urls], columns=['URL'])

            bucket = self.client.get_bucket(self.bucket_name)
            bucket.blob(f'{outputFile}').upload_from_string(df.to_csv(index=False), 'text/csv')
    
    
    # End function
    p = beam.Pipeline(options=options)
    (p | 'Read Input CSV' >> beam.io.ReadFromText(inputFile, skip_header_lines=1)
       | 'Map Codes' >> beam.Map(genURLs)
       | 'Store Output File' >> beam.ParDo(WriteCSVFIle(bucket)))

    p.run()

代码在我的存储桶中生成URL_test_codes.csv,但该文件仅包含一行(不包括“URL”标头),这告诉我我的代码在处理每一行时正在写入/覆盖文件。有没有办法批量写入整个文件的内容,而不是发出一系列更新文件的请求?我是 Python/Dataflow 的新手,非常感谢任何帮助。

【问题讨论】:

    标签: python csv google-cloud-platform google-cloud-storage google-cloud-dataflow


    【解决方案1】:

    让我们指出问题:明显的问题是 GCS 方面的配额问题,反映在“429”错误代码上。但正如您所指出的,这源于固有问题,这与您尝试将数据写入 blob 的方式更相关。

    由于 Beam 管道会生成元素的并行集合,因此当您将元素添加到 PCollection 时,将为每个元素执行每个管道步骤,换句话说,您的 ParDo 函数将尝试每次向您的输出文件写入一些内容PCollection 中的元素。

    因此,您的 WriteCSVFIle 函数存在一些问题。例如,为了将您的 PCollection 写入 GCS,最好使用专注于编写整个 PCollection 的单独管道任务,如下所示:

    首先,您可以导入这个已经包含在 Apache Beam 中的函数:

    from apache_beam.io import WriteToText
    

    然后,您在管道的末端使用它:

    | 'Write PCollection to Bucket' >> WriteToText('gs://{0}/{1}'.format(bucket_name, outputFile))
    

    使用此选项,您无需创建存储客户端或引用 blob,该函数只需要接收 GCS URI,它将写入最终结果,您可以根据找到的参数进行调整在documentation。

    有了这个,您只需要处理在您的 WriteCSVFIle 函数中创建的数据框。每个管道步骤都会创建一个新的 PCollection,因此如果 Dataframe-creator 函数应该从 URL 的 PCollection 接收元素,那么根据您当前的逻辑,从 Dataframe 函数产生的新 PCollection 元素将在每个 url 有 1 个数据帧,但因为它看起来考虑到“URL”是数据框中的唯一列,您只想从 genURLs 写入结果,也许直接从 genURLs 到 WriteToText 可以输出您要查找的内容。

    无论哪种方式,您都可以相应地调整您的管道,但至少通过 WriteToText 转换,它会负责将您的整个最终 PCollection 写入您的 Cloud Storage 存储桶。

    【讨论】:

    • 感谢您的洞察力和建议 - 我删除了 beam.Map() 步骤并用 WriteToText() 替换它,它可以生成包含所有 250k URL 的文件。我现在唯一好奇的是如何更新新文件的元数据,使其 Content-Type 为 text/csv 而不是默认的 text/plain。我看到传入 gsutil 命令可能会完成此操作,但也许有更好/更简单的方法
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-11
    • 1970-01-01
    • 2012-05-27
    • 1970-01-01
    • 1970-01-01
    • 2017-07-12
    相关资源
    最近更新 更多