Mostrando entradas con la etiqueta gcp. Mostrar todas las entradas
Mostrando entradas con la etiqueta gcp. Mostrar todas las entradas

sábado, 9 de mayo de 2026

GCP: optimizando consultas y sentencias en BigQuery


BigQuery nos permite realizar las siguientes operaciones:

  1. Consulta de datos: Puedes ejecutar consultas SQL complejas para extraer datos de tus conjuntos de datos y realizar análisis avanzados.
  2. Análisis de datos en tiempo real: BigQuery admite consultas en tiempo real sobre datos de streaming, lo que te permite analizar y visualizar datos en tiempo real a medida que llegan.
  3. Análisis geoespacial: BigQuery incluye funciones y operaciones para realizar análisis geoespaciales, como cálculos de distancia, intersecciones espaciales y agrupaciones geográficas.
  4. Operaciones de agregación: Puedes realizar operaciones de agregación como SUM, AVG, COUNT, MAX y MIN en tus datos para resumir la información y obtener insights.
  5. Procesamiento de texto: BigQuery proporciona funciones y operadores para procesar datos de texto, como búsquedas de patrones, análisis de sentimientos y extracción de entidades.
  6. Integración con herramientas de análisis y visualización: Puedes integrar BigQuery con herramientas de análisis y visualización de datos populares como Google Data Studio, Tableau y Power BI para crear paneles interactivos y visualizaciones de datos.
  7. Machine Learning: BigQuery ML te permite construir y entrenar modelos de aprendizaje automático directamente en tus datos almacenados en BigQuery, sin necesidad de moverlos a otro lugar.
  8. Carga y exportación de datos: Puedes cargar datos en BigQuery desde archivos locales, Google Cloud Storage, servicios de streaming como Pub/Sub y otras fuentes de datos. También puedes exportar datos desde BigQuery a diferentes formatos de archivo y servicios de almacenamiento.
  9. Seguridad y control de acceso: BigQuery ofrece controles de acceso granulares y opciones de cifrado para proteger tus datos y garantizar la conformidad con las normativas de privacidad.
  10. Administración y monitoreo: BigQuery proporciona herramientas para administrar y monitorear tus recursos, consultas y cargas de trabajo, como el tablero de control de BigQuery y Cloud Monitoring.

Sin embargo, también tiene ciertos limitantes como lo pueden ser:

  • Debes ser cuidadoso por el uso, ya que las cuotas pueden restringir ciertas operaciones iterativas. 
  • Si operas sobre una misma tabla, puede haber bloqueos (lo que evitará un buen almacenamiento de tu información). 
  • No permite subconsultas como en Informix o herramientas similares. 
  • El tamaño de una fila no puede superar los 10MB.
  • Etc.

Y es ahí donde entran las optimizaciones con CTEs (Common Table Expressions) o WITH en las consultas. Imaginemos el siguiente bloque de código:

/*
   Este bloque es para actualizar la información de
 la tabla2 desde la tabla1.
*/

begin
declare fecha_origen date default '2026-04-13';
declare fecha_actual date default '2026-05-09';

for record 
in(select user, info_process from
 `mydataset.tabla1` 
where date_process = fecha_origen) 
do
  update `mydataset.tabla2` 
set user = record.user, 
info_process = record.info_process, 
date_process= fecha_actual 
where date_process = fecha_origen ;
end for;

end;

El bloque realiza la operación de actualización correctamente, pero no es lo más óptimo. Pues debemos cuidar los recursos.

Rehacemos el bloque pero usando la cláusula de WITH. Esto nos permitirá la optimización del bloque y ahorraremos tiempo y recursos valiosos.

Tenemos entonces lo siguiente:

/*
   Este bloque es para actualizar la información de
 la tabla2 desde la tabla1.
*/

begin
declare fecha_origen date default '2026-04-13';
declare fecha_actual date default '2026-05-09';

update `mydataset.tabla2` as tab_up 
set user = src.user, 
info_process = src.info_process, 
date_process= fecha_actual 
from(
   select id, user, info_process 
   FROM `mydataset.tabla1` where date_process = fecha_origen 
) as src where  tab_up.id = src.id 
 and  tab_up.date_process = fecha_origen ;

end;

Aunque no empleamos la cláusula, seguimos su misma lógica: optimizar la consulta.

¿Qué pasaría si quisieramos hacer una inserción?

Tomando en cuenta el siguiente bloque:

/*
   Este bloque es para actualizar la información de
 la tabla2 desde la tabla1.
*/

begin
declare fecha_actual date default '2026-05-09';

for record 
in(select valor_mensual from
 `mydataset.tabla1` 
where date_process = fecha_actual and importe < 99.9) 
do
  
insert into `mydataset.tabla2`(valor_mensual) value (record.valor_mensual); 

end for;

end;

El bloque trabaja casi perfectamente, pero no es óptimo el uso de recursos.

Rehacemos el bloque con la lógica de WITH:

/*
   Este bloque es para actualizar la información de
 la tabla2 desde la tabla1.
*/

begin
declare fecha_actual date default '2026-05-09';


insert into `mydataset.tabla2`(valor_mensual) 
with source as(
  select t.valor_mensual from `mydataset.tabla1` t
  where t.date_process = fecha_actual and t.importe < 99.9
) select * from source;



end;

Como se puede observar si usamos la cláusula WITH y no solo su lógica como en el ejemplo del bloque UPDATE.

Y cómo es lógico, también lo podemos aplicar a las consultas con SELECT. Miremos una consulta no optimizada y comparémosla con una que sí lo está:

declare fecha_actual date default '2026-05-09';
declare importe_max float64;
set importe_max = 99.9;

select user, date_process, importe, valor_mensual
 `mydataset.tabla2`
where date_process = (
   select date_process
 `mydataset.tabla1`where date_process = fecha_actual 
) and importe = (
  select importe
 `mydataset.tabla1`where importe < importe_max 
);

Optimizada:

DECLARE fecha_actual DATE DEFAULT DATE '2026-05-09';
DECLARE importe_max FLOAT64 DEFAULT 99.9;

WITH filtro_fecha AS (
  SELECT date_process
  FROM `mydataset.tabla1`
  WHERE date_process = fecha_actual
),
filtro_importe AS (
  SELECT importe
  FROM `mydataset.tabla1`
  WHERE importe < importe_max
)
SELECT user, date_process, importe, valor_mensual
FROM `mydataset.tabla2`
WHERE date_process IN (SELECT date_process FROM filtro_fecha)
  AND importe IN (SELECT importe FROM filtro_importe);

¿Y qué de las operaciones de borrado?

Sin optimizar:

begin
declare fecha_actual date default '2026-05-09';

for record 
in(select valor_mensual from
 `mydataset.tabla1` 
where date_process = fecha_actual and importe < 99.9) 
do
  
delete from `mydataset.tabla2`where valor_mensual = record.valor_mensual
 and date_process = fecha_actual;

end for;

end;

Optimizada:

DECLARE fecha_actual DATE DEFAULT DATE '2026-05-09';

WITH valores_a_borrar AS (
  SELECT valor_mensual
  FROM `mydataset.tabla1`
  WHERE date_process = fecha_actual
    AND importe < 99.9
)
DELETE FROM `mydataset.tabla2`
WHERE date_process = fecha_actual
  AND valor_mensual IN (SELECT valor_mensual FROM valores_a_borrar);

Como hemos visto, el uso de la cláusula WITH nos permite evitar subconsultas repetitivas y hace el código más legible y eficiente.

Seguiremos hablando de este tema en próximas entregas.

Enlaces:

https://codemonkeyjunior.blogspot.com/2024/04/gcp-google-cloud-bigquery.html
WITH statements in BigQuery SQL (Youtube)

viernes, 29 de agosto de 2025

GCP: manejo de excepciones en BigQuery

 

El manejo de excepciones es un mecanismo para identificar, tratar y recuperar el flujo normal de un programa en ejecución. El lenguaje de BigQuery nos permite gestionar los errores cuando estos ocurren  en nuestras consultas SQL. Miremos unos ejemplos.

El siguiente bloque de BigQuery tratará de ejecutar una consulta a una tabla que no existe (tabla_no_existente). El sub-bloque tiene un manejo de excepción de tipo ``ERROR``.

DECLARE error STRING;
BEGIN
  BEGIN
     -- Primera consulta, fallará.
     SELECT * FROM mydataset.tabla_no_existente;
  EXCEPTION WHEN ERROR THEN
      SET error = 'Ha ocurrido una excepcion.';
      RAISE; -- Termina la ejecución de todo el bloque. La segunda consulta no se realizará.
   END;

   -- Segunda consulta.
   SELECT * FROM mydataset.tabla_existente LIMIT 1;
END;

Si queremos que la segunda consulta se ejecute aún si ocurre un error en la primera consulta, entonces tan solo quitamos la instrucción ``RAISE``, la cual se utiliza para generar explícitamente un error o reactivar una excepción existente.

Tendríamos lo siguiente:

DECLARE error STRING;
BEGIN
  BEGIN
     -- Primera consulta, fallará.
     SELECT * FROM mydataset.tabla_no_existente;
  EXCEPTION WHEN ERROR THEN
      SET error = 'Ha ocurrido una excepcion.';
   END;

   -- Segunda consulta. Se ejecuta aunque la primera haya fallado.
   SELECT * FROM mydataset.tabla_existente LIMIT 1;
END;

La primera consulta fallará, pero no se terminará la ejecución del bloque completo. Se ejecutará la segunda consulta.

¿Cómo obtener información detallada del error ocurrido?

Crearemos un Stored Procedure que nos permitirá mostrar a detalle el error en la ejecución de una consulta SQL en BigQuery.

CREATE OR REPLACE PROCEDURE `mydataset.exec_sp_error`(OUT error_flag STRING)
BEGIN
   -- Esta consulta fallará.
   SELECT usuario,* FROM mydataset.tabla_no_existente;

  EXCEPTION WHEN ERROR THEN
   SET error_flag = "Error: ";
   SET error_flag = concat(salida, @@error.message);
   SET error_flag = concat(salida, ". Causa: "); 
   SET error_flag = concat(salida, @@error.statement_text);  
END;

Invocamos el Stored Procedure y consultamos la variable ``error_flag``:

DECLARE error_flag STRING;
-- Invocamos el Stored Procedure
CALL`mydataset.exec_sp_error`(error_flag);
-- Mostrará el error a detalle
SELECT error_flag;

En lenguajes de programación como Java, C++ y C# se puede hacer algo como esto:

try{
   int divide = 1/0;
}catch(ArithmeticException ex){
   System.err.printf("%s\n", ex.getCause());
   System.err.printf("%s\n", ex.getMessage());
   ex.printStackTrace();
}

Lo cual disparará una excepción de tipo ``ArithmeticException``, pues no podemos dividir un número por cero.

En BigQuery el código equivalente sería:

BEGIN
  SELECT 1/0; -- División por cero. 
EXCEPTION WHEN ERROR THEN
  SELECT @@error.message AS error_message;
  SELECT @@error.stack_trace AS stack_trace;
  SELECT @@error.statement_text AS text_error;
END;

Cuando se lance la excepción se mostrará los errores en la consola de BigQuery.

Continuaremos sobre esta serie de BigQuery más adelante.

Enlaces:

https://codemonkeyjunior.blogspot.com/2024/11/stored-procedures-oracle-plsql-gcp.html
https://cloud.google.com/bigquery/docs/error-messages

domingo, 13 de julio de 2025

GCP Pub/Sub: un servicio de mensajería asíncrona y escalable

Según la documentación oficial, Pub/Sub es:

Un servicio de mensajería asíncrona y escalable que separa los servicios que producen mensajes de aquellos que procesan esos mensajes.

Lo que nos permite hacer es:

Permitir que los servicios se comuniquen de forma asíncrona, con latencias de alrededor de 100 milisegundos.

Se usa para:

Las canalizaciones de integración de datos y estadísticas de transmisión a fin de cargar y distribuir datos. Es igual de efectivo que el middleware orientado a la mensajería para la integración de servicios o como una cola con el fin de paralelizar las tareas.

Pub/Sub te permite crear sistemas de productores y consumidores de eventos, llamados publicadores y suscriptores. Los publicadores se comunican con los suscriptores de forma asíncrona mediante la transmisión de eventos, en lugar de llamadas de procedimiento remoto (RPC) síncronas.

Los publicadores envían eventos al servicio de Pub/Sub, sin importar cómo o cuándo se procesarán estos eventos. Luego, Pub/Sub entrega eventos a todos los servicios que reaccionan a ellos. En los sistemas que se comunican a través de RPC, los publicadores deben esperar a que los suscriptores reciban los datos.

Sin embargo, la integración asíncrona en Pub/Sub aumenta la flexibilidad y solidez del sistema general.

En pocas palabras, es una tecnología similar a RabbitMQ, Apache Kafka y/o ActiveMQ por mencionar solo algunas.

Casos habituales

  • Transferir la interacción del usuario y los eventos del servidor. 
  • Distribución de eventos en tiempo real. 
  • Replicar datos entre bases de datos. 
  • Procesamiento y flujos de trabajo paralelos. 
  • Bus de eventos empresariales. 
  • Transmisión de datos desde aplicaciones, servicios o dispositivos de la IoT. 
  • Actualización de cachés distribuidas. 
  • Balanceo de cargas para la confiabilidad.

Un ejemplo sencillo de esta tecnología sería crear una aplicación que:

  • Publique mensajes a un tópico de Pub/Sub y 
  • Reciba mensajes desde ese tópico.

Para ello es necesario tener una cuenta de GCP. Tener el JDK más actual.

Creamos una aplicación de Spring Boot y en nuestro pom.xml agregamos la siguiente dependencia:

<dependency>
  <groupId>com.google.cloud</groupId>
  <artifactId>spring-cloud-gcp-starter-pubsub</artifactId>
</dependency>

La configuración del archivo YAML sería algo como esto:

# En application.yml
spring:
  cloud:
    gcp:
      pubsub:
        project-id: tu-id-de-proyecto
        credentials:
          location: classpath:tu-archivo-de-credenciales.json

Crearemos una clase tipo Service que servirá como el "publicador":

@Service
public class PublisherService {
    private final PubSubTemplate pubSubTemplate;

    public PublisherService(PubSubTemplate pubSubTemplate) {
        this.pubSubTemplate = pubSubTemplate;
    }

    public void enviarMensaje(String mensaje) {
        pubSubTemplate.publish("mi-topico", mensaje);
    }
}

Crearemos una clase tipo Component que servirá para escuchar los mensajes:

@Component
public class Subscriber {
    @PubSubSubscriber(subscription = "mi-suscripcion")
    public void recibirMensaje(String mensaje) {
        System.out.println("Mensaje recibido: " + mensaje);
    }
}

Ahora crearemos una clase tipo Controller para invocar el método de envío de mensajes:

@RestController
@RequestMapping("/pubsub")
public class PubSubController {

    private final PublisherService publisherService;

    public PubSubController(PublisherService publisherService) {
        this.publisherService = publisherService;
    }

    @PostMapping("/enviar")
    public ResponseEntity<String> enviar(@RequestBody String mensaje) {
        publisherService.enviarMensaje(mensaje);
        return ResponseEntity.ok("Mensaje publicado: " + mensaje);
    }
}

Usando curl para probar el mensaje:

curl -X POST http://localhost:8080/pubsub/enviar \
     -H "Content-Type: text/plain" \
     -d "Hola desde el controlador!"

También agregando un nuevo método en el controller:

@PostMapping("/enviar-json")
public ResponseEntity<String> enviarJson(@RequestBody Map<String, String> payload) {
    String mensaje = payload.get("mensaje");
    publisherService.enviarMensaje(mensaje);
    return ResponseEntity.ok("Mensaje JSON publicado: " + mensaje);
}

Y usando curl:

curl -X POST http://localhost:8080/pubsub/enviar-json \
     -H "Content-Type: application/json" \
     -d '{"mensaje": "¡Hola JSON!"}'

Continuaremos sobre temas de GCP en próximas entregas.

Enlaces:

https://cloud.google.com/pubsub/docs/overview
https://www.rabbitmq.com/
https://kafka.apache.org/
https://activemq-apache-org.translate.goog/

sábado, 28 de junio de 2025

Subir un archivo CSV a una tabla en GCP BigQuery

En una entrega pasada vimos como cargar datos desde un CSV a una tabla en BigQuery.

Esta vez crearemos un programa con Java y BigQuery. Para ello necesitamos tener un archivo con datos el cual llamaremos "DATOS.txt". El contenido del archivo será similar a esto:

20250412,34,450.0
20250412,34,432.0
20250412,122,500.0

Tendremos una tabla llamada ``tkdata`` con los siguientes campos:

fchinf DATE 
numregs INT64 
valor STRING 

Requisitos:

  • Tener acceso a GCP BigQuery.
  • Tener JDK 11 o más actual. 
  • Tener Maven (el más actual).

¿Qué hará el programa?

  1. Verificar si existe el archivo a subir ("DATOS.txt"). 
  2. Obtener su contenido y guardarlo en una lista tipo String. 
  3. Crear un objeto tipo StringBuilder a partir de la lista tipo String. 
  4. Enviar el contenido del objeto StringBuilder a un nuevo archivo ("tkdata.csv") que se guardará en el mismo bucket del archivo original. 
  5. Cargar el contenido del nuevo archivo a la tabla ``tkdata``.
  6. Verificar que los datos hayan sido cargados a la tabla.

BigQueryCsvUploader.java

import com.google.cloud.bigquery.*;
import com.google.cloud.storage.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.BufferedReader;
import java.io.IOException;
import java.io.StringReader;
import java.nio.charset.StandardCharsets;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;

public class BigQueryCsvUploader {
    private static final Logger LOG = LoggerFactory.getLogger(BigQueryCsvUploader.class);
    private static final String NAME_FILE_ORIGINAL = "DATOS.csv";
    private static final String TABLE_NAME = "tkdata";
    private static final String DATASET = "mydataset";
    private static final String PROJECT = "myproject";
    private static final String BUCKET = "mybucket";
    private static final String NAME_NEW_FILE = "tkdata.csv";
    private static final long MAX_SIZE_BYTES = 300 * 1024 * 1024; // 300 MB

   
    public static class Tkdata {
        private Date fchinf;
        private long numregs;
        private String valor;

        public Date getFchinf() { return fchinf; }
        public void setFchinf(String fchinf) throws ParseException {
            SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
            this.fchinf = sdf.parse(fchinf);
        }
        public long getNumregs() { return numregs; }
        public void setNumregs(String numregs) { this.numregs = Long.parseLong(numregs); }
        public String getValor() { return valor; }
        public void setValor(String valor) { this.valor = valor; }
    }

   
    public static boolean existFile(String bucketName, String fileName) {
        Storage storage = StorageOptions.getDefaultInstance().getService();
        Blob blob = storage.get(BlobId.of(bucketName, fileName));
        return blob != null && blob.exists();
    }

   
    public static boolean maxSize(String bucketName, 
     String fileName, long maxSizeBytes) {
        Storage storage = StorageOptions.getDefaultInstance().getService();
        Blob blob = storage.get(BlobId.of(bucketName, fileName));
        if (blob == null) {
            LOG.error("El archivo {} no existe en el bucket {}", fileName, bucketName);
            return false;
        }
        return blob.getSize() <= maxSizeBytes;
    }

    
    public static List<Tkdata> getListTkdata(String bucketName, String fileName) {
        List<Tkdata> listaTkdata = new ArrayList<>();
        Storage storage = StorageOptions.getDefaultInstance().getService();
        Blob blob = storage.get(BlobId.of(bucketName, fileName));
        if (blob == null) {
            LOG.error("Archivo{} no encontrado en el bucket {}", fileName, bucketName);
            return listaTkdata;
        }

        String content = new String(blob.getContent(), StandardCharsets.UTF_8);
        try (BufferedReader reader = new BufferedReader(new StringReader(content))) {
            String line;
            reader.readLine();
            while ((line = reader.readLine()) != null) {
                String[] parts = line.split(",", -1); 
                if (parts.length == 3) {
                    Tkdata obj = new Tkdata();
                    try {
                        obj.setFchinf(parts[0].trim());
                        obj.setNumregs(parts[1].trim());
                        obj.setValor(parts[2].trim());
                        listaTkdata.add(obj);
                    } catch (ParseException | NumberFormatException e) {
                        LOG.error("Error parseando linea: {}", line, e);
                    }
                } else {
                    LOG.warn("Linea malformada: {}", line);
                }
            }
        } catch (IOException e) {
            LOG.error("Error al leer el contenido", e);
        }
        return listaTkdata;
    }

    
    public static StringBuilder convertToStringBuilder(List<Tkdata> listaTkdata) {
        StringBuilder sb = new StringBuilder();
        SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
        for (Tkdata item : listaTkdata) {
            sb.append(sdf.format(item.getFchinf())).append(",");
            sb.append(item.getNumregs()).append(",");
            sb.append(item.getValor()).append("\n");
        }
        return sb;
    }

   
    public static boolean creaNuevoFileCSV(String bucketName, 
       String fileName, String content) {
        try {
            Storage storage = StorageOptions.getDefaultInstance().getService();
            BlobId blobId = BlobId.of(bucketName, fileName);
            BlobInfo blobInfo = BlobInfo.newBuilder(blobId).setContentType("text/csv").build();
            storage.create(blobInfo, content.getBytes(StandardCharsets.UTF_8));
            return true;
        } catch (Exception e) {
            LOG.error("Error al subir archivo {} al bucket {}", fileName, bucketName, e);
            return false;
        }
    }

    
    public static boolean loadCSVToTableBigQuery(String projectId, 
     String datasetId, String bucketName, String sourceFileName, String tableName) {
        try {
            BigQuery bigquery = BigQueryOptions.getDefaultInstance().getService();
            TableId tableId = TableId.of(projectId, datasetId, tableName);

            
            Schema schema = Schema.of(
                Field.of("fchinf", StandardSQLTypeName.DATE),
                Field.of("numregs", StandardSQLTypeName.INT64),
                Field.of("valor", StandardSQLTypeName.STRING)
            );

           
            JobConfiguration jobConfig = LoadJobConfiguration.newBuilder(
                tableId,
                String.format("gs://%s/%s", bucketName, sourceFileName),
                FormatOptions.csv()
            )
                .setSchema(schema)
                .setSkipLeadingRows(0) 
                .setWriteDisposition(JobInfo.WriteDisposition.WRITE_APPEND) 
                .build();

            
            Job job = bigquery.create(JobInfo.of(jobConfig));
            job = job.waitFor();
            if (job.isDone() && job.getStatus().getError() == null) {
                LOG.info("CSV {} successfully loaded into BigQuery table {}.{}", 
                 sourceFileName, datasetId, tableName);
                return true;
            } else {
                LOG.error("Error cargando CSV a BigQuery: {}", job.getStatus().getError());
                return false;
            }
        } catch (Exception e) {
            LOG.error("Error en el proceso de carga", e);
            return false;
        }
    }

  
    public static boolean validateTableData(String projectId, 
      String datasetId, String tableName, List<Tkdata> expectedData) {
        try {
            BigQuery bigquery = BigQueryOptions.getDefaultInstance().getService();
            String query = String.format("SELECT fchinf, numregs, valor FROM %s.%s.%s", 
                        projectId, datasetId, tableName);
            QueryJobConfiguration queryConfig = QueryJobConfiguration.newBuilder(query).build();

           
            TableResult result = bigquery.query(queryConfig);
            long rowCount = result.getTotalRows();

            if (rowCount == 0) {
                LOG.error("No hay datos en la tabla {}.{}", datasetId, tableName);
                return false;
            }

           
            if (expectedData != null && rowCount != expectedData.size()) {
                LOG.warn("Error en conteo de datos {}, found {}", expectedData.size(), rowCount);
                return false;
            }

            
            LOG.info("Sample data from {}.{} (first 5 rows):", datasetId, tableName);
            SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
            int maxRowsToLog = 5;
            int rowIndex = 0;
            for (FieldValueList row : result.iterateAll()) {
                if (rowIndex >= maxRowsToLog) break;
                String fchinf = row.get("fchinf").getStringValue(); 
                long numregs = row.get("numregs").getLongValue();
                String valor = row.get("valor").getStringValue();
                LOG.info("Row {}: fchinf={}, numregs={}, valor={}", rowIndex + 1, fchinf, numregs, valor);
                rowIndex++;
            }

            LOG.info("Validacion exitosa: {} rows found in table {}.{}", 
          rowCount, datasetId, tableName);
            return true;
        } catch (Exception e) {
            LOG.error("Error validando datos en la tabla {}.{}", datasetId, tableName, e);
            return false;
        }
    }

    public static void main(String[] args) {
     
        if (!existFile(BUCKET, NAME_FILE_ORIGINAL)) {
            LOG.error("Archivo{} no existe en el buckett {}.", NAME_FILE_ORIGINAL, BUCKET);
            return;
        }

        
        if (!maxSize(BUCKET, NAME_FILE_ORIGINAL, MAX_SIZE_BYTES)) {
            LOG.error("El archivo {} excede el tamaño: 300MB.", NAME_FILE_ORIGINAL);
            return;
        }

        
        List<Tkdata> listaTkdata = getListTkdata(BUCKET, NAME_FILE_ORIGINAL);
        if (listaTkdata.isEmpty()) {
            LOG.error("No hay datos para leer {}. ", NAME_FILE_ORIGINAL);
            return;
        }

        StringBuilder sb = convertToStringBuilder(listaTkdata);

        
        if (creaNuevoFileCSV(BUCKET, NAME_NEW_FILE, sb.toString())) {
            LOG.info("Archivo {} cargado al bucket {}", NAME_NEW_FILE, BUCKET);
            LOG.info("Carga de datos {} a la tabla: {}.{}", NAME_NEW_FILE, DATASET, TABLE_NAME);
            if (loadCSVToTableBigQuery(PROJECT, DATASET, BUCKET, NAME_NEW_FILE, TABLE_NAME)) {
                if (validateTableData(PROJECT, DATASET, TABLE_NAME, listaTkdata)) {
                    LOG.info("Validacion correcta {}.{}", DATASET, TABLE_NAME);
                } else {
                    LOG.error("Validacion fallida {}.{}", DATASET, TABLE_NAME);
                }
            } else {
                LOG.error("Fallo al cargar los datos a la tabla de BigQuery  {}.{}", DATASET, TABLE_NAME);
            }
        } else {
            LOG.error("Fallo al cargar el archivo {} al bucket {}", NAME_NEW_FILE, BUCKET);
        }
    }
}

Es importante hacer notar que los tipos de datos en el archivo deben concordar a los datos de la tabla. En caso contrario, no cargarán.

También es importante tener estas dependencias en el pom.xml

<dependencies>
    <dependency>
        <groupId>com.google.cloud</groupId>
        <artifactId>google-cloud-storage</artifactId>
        <version>2.44.0</version>
    </dependency>
    <dependency>
        <groupId>com.google.cloud</groupId>
        <artifactId>google-cloud-bigquery</artifactId>
        <version>2.44.0</version>
    </dependency>
    <dependency>
        <groupId>org.slf4j</groupId>
        <artifactId>slf4j-api</artifactId>
        <version>2.0.13</version>
    </dependency>
    <dependency>
        <groupId>org.slf4j</groupId>
        <artifactId>slf4j-simple</artifactId>
        <version>2.0.13</version>
    </dependency>
</dependencies>

De preferncia usar Eclipse IDE para generar un .jar o ejecutarlo directamente.

Si la ejecución es correcta podremos ver los datos cargados en la tabla. Continuaremos más sobre BigQuery en próximas entregas.

Enlaces:

https://cloud.google.com/bigquery/docs/loading-data-cloud-storage-csv

Jai un lenguaje de programación inspirado en C++

Hoy hablaremos de un nuevo lenguaje de programación llamado Jai . Se trata de un lenguaje de programación que está desarrolland...

Etiquetas

Archivo del blog