A task interesting enough for me to see one of the practical applications of concurrency. The mission was to backfill device data into Bigtable spanning the past x years. For context, the end goal is to have proper management of devices; there are too many old (read: ancient) devices in the database and it’s just been piling up, so we needed a better way to track active devices.

I initially thought this would take around 2-3 days to run. It would’ve been great to make use of Bigtable’s batch write since this is considered a large scale writes with millions of rows. However, this doesn’t seem to fit into my use case because of the conditional aspect – reading each row and updating it only if LastSeen has not been updated.

It is possible to spin up more pods too but this is probably an overkill for this task and I’ll have to deal with duplication in processing the rows. Leveraging on Goroutine seemed like the surefire (and most convenient) way to get this done faster.

With the added Goroutine concurrency, the write rate seems to increase proportionally to the number of Goroutines.

Le logo de Jekyll

Though this is not the case for CPU usage. On average, it uses 18% CPU with no concurrency and 66% with 5 Goroutines. This is kinda expected since the workload is mainly I/O-bound operations.

Le logo de Jekyll

In the end, the job took 4 hours to finish with 20 Goroutines.

A short code snippet for reference.

counter := 0
const numWorkers = 5
var wg sync.WaitGroup
var mu sync.Mutex

for i := 0; i < numWorkers; i++ {
  wg.Add(1)
  go func() {
    defer wg.Done()
    for row := range rows {
      structuredRow := StructuredRow{
        DeviceID: row[0],
        UserID:   row[1],
        LastSeen: getLastSeen(row),
      }

      // Write rows based on filter
      var columns []bigtable.BigtableColumn
      columns = append(columns, bigtable.BigtableColumn{
        ColumnFamily:    "last_seen",
        ColumnQualifier: "time",
        Timestamp:       time.Now().UnixMicro(),
        Value:           []byte(strconv.Itoa(structuredRow.LastSeen)),
      })

      rowKey := fmt.Sprintf("%s:%s", structuredRow.DeviceID, structuredRow.UserID)
      btRow := bigtable.BigtableRow{
        RowKey:  rowKey,
        Columns: columns,
      }

      err = BigtableClient.ConditionalWriteRow(ctx, btRow, structuredRow.LastSeen)
      if err != nil {
        fmt.Println("Error writing to bigtable:", err)
      } else {
        mu.Lock()
        fmt.Printf("Worker[%d] Row=%d --- Written to bigtable: %s, %d\n", i, counter, rowKey, structuredRow.LastSeen)
        counter++
        mu.Unlock()
      }
    }
  }()
}

wg.Wait()