Empresas
Empregos
  • Sobre nós
  • Soluções
    • Publicação de vagas
      Publique sua vaga e receba candidatos qualificados em 48h.
    • Avaliações de candidatos
      Mais de 500 testes técnicos e psicológicos, mais anti-fraude.
    • Headhunting
      Busca executiva personalizada do início ao fim.
    • Folha de Pagamento + EOR
      Dispersão de folha e EOR em mais de 15 países da LATAM.
  • Preços
  • Empregos

0

405
Visualizações
optimización de costos de trabajo de flujo de datos de google

Ejecuté el siguiente código para 522 archivos gzip de 100 GB de tamaño y, después de descomprimirlos, serán alrededor de 320 GB de datos y datos en formato protobuf y escribiré la salida en GCS. He usado n1 máquinas y regiones estándar para la entrada, la salida se ha cuidado y el trabajo me costó alrededor de 17 $, esto es para datos de media hora, por lo que realmente necesito hacer una optimización de costos aquí muy desesperadamente.

Costo que obtengo de la siguiente consulta

 SELECT l.value AS JobID, ROUND(SUM(cost),3) AS JobCost FROM `PROJECT.gcp_billing_data.gcp_billing_export_v1_{}` bill, UNNEST(bill.labels) l WHERE service.description = 'Cloud Dataflow' and l.key = 'goog-dataflow-job-id' and extract(date from _PARTITIONTIME) > "2020-12-31" GROUP BY 1

Código completo

 import time import sys import argparse import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.options.pipeline_options import SetupOptions import csv import base64 from google.protobuf import timestamp_pb2 from google.protobuf.json_format import MessageToDict from google.protobuf.json_format import MessageToJson import io import logging from io import StringIO from google.cloud import storage import json ###PROTOBUF CLASS from otherfiles import processor_pb2 class ConvertToJson(beam.DoFn): def process(self, message, *args, **kwargs): import base64 from otherfiles import processor_pb2 from google.protobuf.json_format import MessageToDict from google.protobuf.json_format import MessageToJson import json if (len(message) >= 4): b64ProtoData = message[2] totalProcessorBids = int(message[3] if message[3] and message[3] is not None else 0); b64ProtoData = b64ProtoData.replace('_', '/') b64ProtoData = b64ProtoData.replace('*', '=') b64ProtoData = b64ProtoData.replace('-', '+') finalbunary = base64.b64decode(b64ProtoData) log = processor_pb2.ProcessorLogProto() log.ParseFromString(finalbunary) #print(log) jsonObj = MessageToDict(log,preserving_proto_field_name=True) jsonObj["totalProcessorBids"] = totalProcessorBids #wjdata = json.dumps(jsonObj) print(jsonObj) return [jsonObj] else: pass class ParseFile(beam.DoFn): def process(self, element, *args, **kwargs): import csv for line in csv.reader([element], quotechar='"', delimiter='\t', quoting=csv.QUOTE_ALL, skipinitialspace=True): #print (line) return [line] def run(): parser = argparse.ArgumentParser() parser.add_argument("--input", dest="input", required=False) parser.add_argument("--output", dest="output", required=False) parser.add_argument("--bucket", dest="bucket", required=True) parser.add_argument("--bfilename", dest="bfilename", required=True) app_args, pipeline_args = parser.parse_known_args() #pipeline_args.extend(['--runner=DirectRunner']) pipeline_options = PipelineOptions(pipeline_args) pipeline_options.view_as(SetupOptions).save_main_session = True bucket_input=app_args.bucket bfilename=app_args.bfilename storage_client = storage.Client() bucket = storage_client.get_bucket(bucket_input) blob = bucket.blob(bfilename) blob = blob.download_as_string() blob = blob.decode('utf-8') blob = StringIO(blob) pqueue = [] names = csv.reader(blob) for i,filename in enumerate(names): if filename and filename[0]: pqueue.append(filename[0]) with beam.Pipeline(options=pipeline_options) as p: if(len(pqueue)>0): input_list=app_args.input output_list=app_args.output events = ( p | "create PCol from list" >> beam.Create(pqueue) | "read files" >> beam.io.textio.ReadAllFromText() | "Transform" >> beam.ParDo(ParseFile()) | "Convert To JSON" >> beam.ParDo(ConvertToJson()) | "Write to BQ" >> beam.io.WriteToBigQuery( table='TABLE', dataset='DATASET', project='PROJECT', schema="dataevent:STRING", create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, insert_retry_strategy=RetryStrategy.RETRY_ON_TRANSIENT_ERROR, custom_gcs_temp_location='gs://BUCKET/gcs-temp-to-bq/', method='FILE_LOADS')) ##bigquery failed rows NOT WORKING so commented #(events[beam.io.gcp.bigquery.BigQueryWriteFn.FAILED_ROWS] | "Bad lines" >> beam.io.textio.WriteToText("error_log.txt")) ##WRITING TO GCS #printFileConetent | "Write TExt" >> beam.io.WriteToText(output_list+"file_",file_name_suffix=".json",num_shards=1, append_trailing_newlines = True) if __name__ == '__main__': logging.getLogger().setLevel(logging.INFO) run()

El trabajo tomó alrededor de 49 minutos.

Cosas que probé: 1) Para avro, generé un esquema que debe estar en JSON para el archivo proto e intenté debajo del código para convertir un diccionario a avro msg, pero lleva tiempo ya que el tamaño del diccionario es mayor. schema_separated= es un esquema avro JSON y funciona bien

 with beam.Pipeline(options=pipeline_options) as p: if(len(pqueue)>0): input_list=app_args.input output_list=app_args.output p1 = p | "create PCol from list" >> beam.Create(pqueue) readListofFiles=p1 | "read files" >> beam.io.textio.ReadAllFromText() parsingProtoFile = readListofFiles | "Transform" >> beam.ParDo(ParseFile()) printFileConetent = parsingProtoFile | "Convert To JSON" >> beam.ParDo(ConvertToJson()) compressIdc=True use_fastavro=True printFileConetent | 'write_fastavro' >> WriteToAvro( output_list+"file_", # '/tmp/dataflow/{}/{}'.format( # 'demo', 'output'), # parse_schema(json.loads(SCHEMA_STRING)), parse_schema(schema_separated), use_fastavro=use_fastavro, file_name_suffix='.avro', codec=('deflate' if compressIdc else 'null'), )
  1. En el código principal, traté de insertar el registro JSON como una cadena en la tabla de bigquery para poder usar funciones JSON en bigquery para extraer los datos y eso tampoco salió bien y obtuve el siguiente error.

    mensaje: 'Error al leer los datos, mensaje de error: la tabla JSON encontró demasiados errores, se rindió. Filas: 1; errores: 1. Consulte la colección de errores[] para obtener más detalles.' motivo: 'no válido'> [mientras se ejecuta 'Escribir en BQ/BigQueryBatchFileLoads/WaitForDestinationLoadJobs']

  2. Intenté insertar el diccionario JSON anterior en bigquery proporcionando el esquema JSON a la tabla y también funciona bien

Ahora el desafío es el tamaño después de deserializar el prototipo a JSON dict se duplica y el costo se calculará en el flujo de datos según la cantidad de datos procesados

Estoy tratando y leyendo mucho para hacer que esto funcione y, si funciona, entonces puedo hacerlo estable para la producción.

Ejemplo de registro JSON.

 {'timestamp': '1609286400', 'bidResponseId': '5febc300000115cd054b9fd6840a5af1', 'aggregatorId': '1', 'userId': '7567d74e-2e43-45f4-a42a-8224798bb0dd', 'uniqueResponseId': '', 'adserverId': '1002418', 'dataVersion': '1609285802', 'geoInfo': {'country': '101', 'region': '122', 'city': '11605', 'timezone': '420'}, 'clientInfo': {'os': '4', 'browser': '1', 'remoteIp': '36.70.64.0'}, 'adRequestInfo': {'requestingPage': 'com.opera.mini.native', 'siteId': '557243954', 'foldPosition': '2', 'adSlotId': '1', 'isTest': False, 'opType': 'TYPE_LEARNING', 'mediaType': 'BANNER'}, 'userSegments': [{'id': '2029660', 'weight': -1.0, 'recency': '1052208'}, {'id': '2034588', 'weight': -1.0, 'recency': '-18101'}, {'id': '2029658', 'weight': -1.0, 'recency': '744251'}, {'id': '2031067', 'weight': -1.0, 'recency': '1162398'}, {'id': '2029659', 'weight': -1.0, 'recency': '862833'}, {'id': '2033498', 'weight': -1.0, 'recency': '802749'}, {'id': '2016729', 'weight': -1.0, 'recency': '1620540'}, {'id': '2034584', 'weight': -1.0, 'recency': '111571'}, {'id': '2028182', 'weight': -1.0, 'recency': '744251'}, {'id': '2016726', 'weight': -1.0, 'recency': '1620540'}, {'id': '2028183', 'weight': -1.0, 'recency': '744251'}, {'id': '2028178', 'weight': -1.0, 'recency': '862833'}, {'id': '2016722', 'weight': -1.0, 'recency': '1675814'}, {'id': '2029587', 'weight': -1.0, 'recency': '38160'}, {'id': '2028177', 'weight': -1.0, 'recency': '862833'}, {'id': '2016719', 'weight': -1.0, 'recency': '1675814'}, {'id': '2027404', 'weight': -1.0, 'recency': '139031'}, {'id': '2028172', 'weight': -1.0, 'recency': '1052208'}, {'id': '2028173', 'weight': -1.0, 'recency': '1052208'}, {'id': '2034058', 'weight': -1.0, 'recency': '1191459'}, {'id': '2016712', 'weight': -1.0, 'recency': '1809526'}, {'id': '2030025', 'weight': -1.0, 'recency': '1162401'}, {'id': '2015235', 'weight': -1.0, 'recency': '139031'}, {'id': '2027712', 'weight': -1.0, 'recency': '139031'}, {'id': '2032447', 'weight': -1.0, 'recency': '7313670'}, {'id': '2034815', 'weight': -1.0, 'recency': '586825'}, {'id': '2034811', 'weight': -1.0, 'recency': '659366'}, {'id': '2030004', 'weight': -1.0, 'recency': '139031'}, {'id': '2027316', 'weight': -1.0, 'recency': '1620540'}, {'id': '2033141', 'weight': -1.0, 'recency': '7313670'}, {'id': '2034736', 'weight': -1.0, 'recency': '308252'}, {'id': '2029804', 'weight': -1.0, 'recency': '307938'}, {'id': '2030188', 'weight': -1.0, 'recency': '3591519'}, {'id': '2033449', 'weight': -1.0, 'recency': '1620540'}, {'id': '2029672', 'weight': -1.0, 'recency': '1441083'}, {'id': '2029664', 'weight': -1.0, 'recency': '636630'}], 'perfInfo': {'timeTotal': '2171', 'timeBidInitialize': '0', 'timeProcessDatastore': '0', 'timeGetCandidates': '0', 'timeAdFiltering': '0', 'timeEcpmComputation': '0', 'timeBidComputation': '0', 'timeAdSelection': '0', 'timeBidSubmit': '0', 'timeTFQuery': '0', 'timeVWQuery': '8'}, 'learningPercent': 0.10000000149011612, 'pageLanguageId': '0', 'sspUserId': 'CAESECHFlNeuUm16IYThguoQ8ck_1', 'minEcpm': 0.12999999523162842, 'adSpotId': '1', 'creativeSizes': [{'width': '7', 'height': '7'}], 'pageTypeId': '0', 'numSlots': '0', 'eligibleLIs': [{'type': 'TYPE_OPTIMIZED', 'liIds': [{'id': 44005, 'reason': '12', 'creative_id': 121574, 'bid_amount': 8.403361132251052e-08}, {'id': 46938, 'reason': '12', 'creative_id': 124916, 'bid_amount': 8.403361132251052e-06}, {'id': 54450, 'reason': '12', 'creative_id': 124916, 'bid_amount': 2.0117618771650174e-05}, {'id': 54450, 'reason': '12', 'creative_id': 135726, 'bid_amount': 2.4237295484638312e-05}]}, {'type': 'TYPE_LEARNING'}], 'bidType': 4, 'isSecureRequest': True, 'sourceType': 3, 'deviceBrand': 82, 'deviceModel': 1, 'sellerNetworkId': 12814, 'interstitialRequest': False, 'nativeAdRequest': True, 'native': {'mainImg': [{'w': 0, 'h': 0, 'wmin': 1200, 'hmin': 627}, {'w': 0, 'h': 0, 'wmin': 1200, 'hmin': 627}, {'w': 0, 'h': 0, 'wmin': 1200, 'hmin': 627}, {'w': 0, 'h': 0, 'wmin': 1200, 'hmin': 627}], 'iconImg': [{'w': 0, 'h': 0, 'wmin': 0, 'hmin': 0}, {'w': 0, 'h': 0, 'wmin': 100, 'hmin': 100}, {'w': 0, 'h': 0, 'wmin': 0, 'hmin': 0}, {'w': 0, 'h': 0, 'wmin': 100, 'hmin': 100}], 'logoImg': [{'w': 0, 'h': 0, 'wmin': 100, 'hmin': 100}, {'w': 0, 'h': 0, 'wmin': 0, 'hmin': 0}, {'w': 0, 'h': 0, 'wmin': 100, 'hmin': 100}, {'w': 0, 'h': 0, 'wmin': 0, 'hmin': 0}]}, 'throttleWeight': 1, 'isSegmentReceived': False, 'viewability': 46, 'bannerAdRequest': False, 'videoAdRequest': False, 'mraidAdRequest': True, 'jsonModelCallCount': 0, 'totalProcessorBids': 1}

¿Puede alguien ayudarme aquí?

Capturas de pantalla de PFA también como referencia ingrese la descripción de la imagen aquí

ingrese la descripción de la imagen aquí

over 4 years ago · Hanz Gallego
2 Respostas
Responde à pergunta

0

Mi consejo aquí sería usar Java para realizar sus transformaciones.

En Java, puede convertir el Protobuf en Avro de esta manera: Escribiendo el objeto protobuf en parquet usando apache beam

Y una vez que haya hecho eso, puede usar AvroIO para escribir los datos en los archivos.

Java tiene mucho más rendimiento que Python y le ahorrará recursos informáticos. Dado que este trabajo hace algo muy simple y no requiere ninguna biblioteca especial de Python, lo animo encarecidamente a que intente usar Java.

over 4 years ago · Hanz Gallego Relatório

0

Solo quería llamar su atención sobre " FlexRS " si no ha verificado esto. Esto utiliza instancias de máquinas virtuales (VM) interrumpibles y de esa manera puede reducir su costo.

over 4 years ago · Hanz Gallego Relatório
Responde à pergunta
Encontrar trabalhos remotos

Descubra a nova forma de encontrar um emprego!

melhores empregos
Principais categorias de trabalho
Empresas
Postar vaga Preços Comercial
Jurídico
Termos e Condições Política de privacidade
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomende algumas ofertas para mim
Preciso de ajuda