Slow bulk insert in Spanner

Viewed 483

I'm trying to load a dataset of 1TB in spanner and I can't get further than 5 MB/s with 3 nodes.

Main issue is that the dataset I want to load is mainly composed by integers + some nulls columns, so the commit size is too small, lower than 500K.

I've followed the rules for bulk insert defined at https://cloud.google.com/spanner/docs/bulk-loading. (partitions, workers, etc...) If I add a text column to make the commit size bigger (2.5MB), I can reach a throughput of 70MB/s.

However, with a dataset composed by integer and nulls I don't know how to make the importer faster than 5MB/s.

Table definition:

    id INT64 NOT NULL,
    date DATE NOT NULL,
    type STRING(16) NOT NULL,
    category STRING(3) NOT NULL,
    quadkey STRING(18) NOT NULL,
    subcategory STRING(2) NOT NULL,
    txn INT64,
    accouunts INT64,
    acct_cnt FLOAT64,
    avg_freq FLOAT64,
    avg_spend_amt FLOAT64,
    avg_ticket FLOAT64,
    txn_amt FLOAT64,
    txn_cnt FLOAT64,
) PRIMARY KEY (id)
2 Answers

There are many factors that contribute to a maximized throughput when loading data in Spanner.

1- The first step is to ensure that there isn’t any network capacity or added latencies due to the datasets being loaded into Spanner from an inefficient storage solution. Having your datasets in a Cloud Storage bucket in the same region as your Spanner instance and loading from there is a safe way to do so.

2- Then Spanner needs to partition the data by primary key, so that the splits are equally distributed and the nodes have equal workload to insure best performances. Spanner organizes rows lexicographically, so here are the recommendations on how to optimize the splits. This will avoid “hotspots”, where some nodes have higher workload than others. This can be assessed by the CPU utilization - high priority graph accessible via the console. Trying to get as close to 65% as possible for that metric (the recommended limit) is a good way to insure hotspots are not happening.

3- Having many splits can sometimes seem to increase efficiency, but this benefit is reduced by the coordination overhead created. This is where appropriate ordered batches will find the right balance to maximize throughput. This blog post covers the concept with a detailed example, I highly suggest you read it in full.

As mentioned in the documentation, it is possible to achieve bulk writes of 10-20 MB/s, but since the commits are not in the 1-5 MB at a time, this might not be possible even following the best practices. In this case, the number of operations/s, rows/s or requests/s might be more relevant, and 6MB/s could be in fact a good throughput considering the table schemas. These metrics can be accessed via the console or Stackdriver Monitoring for more advanced insights.

Related