import { BigQuery } from '@google-cloud/bigquery'
import { commonAbis, evmEventDecoder, evmPortalStream } from '@subsquid/pipes/evm'
import { bigqueryTarget } from '@subsquid/pipes/targets/bigquery'
const bigquery = new BigQuery({ projectId: 'my-gcp-project' })
await evmPortalStream({
id: 'erc20-transfers',
portal: 'https://portal.sqd.dev/datasets/ethereum-mainnet',
outputs: evmEventDecoder({
range: { from: '0' },
events: { transfers: commonAbis.erc20.events.Transfer },
}),
}).pipeTo(
bigqueryTarget({
client: { bigquery },
dataset: 'eth_transfers',
tables: [
{
table: 'transfers',
blockNumberColumn: 'block_number',
schema: [
{ name: 'block_number', type: 'INT64', mode: 'REQUIRED' },
{ name: 'log_index', type: 'INT64', mode: 'REQUIRED' },
// TIMESTAMP wire format is INT64 microseconds since epoch
{ name: 'block_timestamp', type: 'TIMESTAMP', mode: 'REQUIRED' },
{ name: 'token', type: 'STRING', mode: 'REQUIRED' },
{ name: 'from', type: 'STRING', mode: 'REQUIRED' },
{ name: 'to', type: 'STRING', mode: 'REQUIRED' },
{ name: 'amount_raw', type: 'STRING', mode: 'REQUIRED' },
],
clusterBy: ['token', 'from'],
},
],
onData: async ({ store, data }) => {
store.insert(
'transfers',
data.transfers.map((t) => ({
block_number: t.block.number,
log_index: t.rawEvent.logIndex,
// The Storage Write API does not parse Date/ISO strings; pass microseconds
block_timestamp: t.timestamp ? t.timestamp.getTime() * 1000 : 0,
token: t.rawEvent.address,
from: t.event.from,
to: t.event.to,
amount_raw: t.event.value.toString(),
})),
)
},
}),
)