You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Google Cloud Function给CSV新增文件名列后导入Google Cloud SQL?

给Cloud SQL导入的CSV数据新增文件名列的解决方案

当然可以实现,这里提供几种适合你现有Node.js Cloud Function场景的方案:

方案一:预处理CSV文件,添加文件名列后再导入

这是最直接的方式,先修改CSV内容,给每行加上文件名列(包括表头),再导入到Cloud SQL。

步骤说明

  1. 从Cloud Storage下载目标CSV到本地临时文件
  2. 读取CSV,给表头新增file_name字段,给每一行数据追加去除后缀后的文件名
  3. 将修改后的CSV上传到Cloud Storage的临时位置(也可在内存流式处理,无需写本地文件)
  4. 调用SQL Admin API导入修改后的CSV
  5. 清理临时资源

代码示例

const { Storage } = require('@google-cloud/storage');
const csv = require('csv-parser');
const createCsvWriter = require('csv-writer').createObjectCsvWriter;
const path = require('path');
const os = require('os');
const fs = require('fs');

const storage = new Storage();
const sqlAdmin = require('@google-cloud/sql-admin').v1beta4;
const client = new sqlAdmin.SqlAdminServiceClient();

async function processAndImportCsv(bucketName, csvFileName, sqlInstance, sqlDatabase, sqlTable) {
  // 1. 下载CSV到临时文件
  const tempDir = os.tmpdir();
  const inputFilePath = path.join(tempDir, csvFileName);
  await storage.bucket(bucketName).file(csvFileName).download({ destination: inputFilePath });

  // 2. 读取并处理CSV
  const rows = [];
  let headers = [];
  const fileNameWithoutExt = path.parse(csvFileName).name;

  await new Promise((resolve, reject) => {
    fs.createReadStream(inputFilePath)
      .pipe(csv())
      .on('headers', (h) => {
        headers = [...h, 'file_name']; // 新增表头
      })
      .on('data', (row) => {
        rows.push({ ...row, file_name: fileNameWithoutExt }); // 追加文件名
      })
      .on('end', resolve)
      .on('error', reject);
  });

  // 3. 写入修改后的CSV到临时文件
  const outputFilePath = path.join(tempDir, `processed_${csvFileName}`);
  const csvWriter = createCsvWriter({
    path: outputFilePath,
    header: headers.map(h => ({ id: h, title: h }))
  });
  await csvWriter.writeRecords(rows);

  // 4. 上传处理后的CSV到Cloud Storage
  const processedFile = await storage.bucket(bucketName).upload(outputFilePath, {
    destination: `processed/${csvFileName}` // 存到processed目录区分原文件
  });

  // 5. 调用SQL Admin API导入处理后的CSV
  const importRequest = {
    project: process.env.GCP_PROJECT,
    instance: sqlInstance,
    resource: {
      importContext: {
        fileType: 'CSV',
        uri: `gs://${bucketName}/processed/${csvFileName}`,
        database: sqlDatabase,
        csvImportOptions: {
          table: sqlTable,
          columns: headers // 指定包含file_name的列
        }
      }
    }
  };

  await client.importInstance(importRequest);

  // 6. 清理临时文件和Cloud Storage的临时文件
  fs.unlinkSync(inputFilePath);
  fs.unlinkSync(outputFilePath);
  await processedFile[0].delete();
}

方案二:导入后批量添加文件名列

如果不想修改原CSV,可以先导入数据到已有file_name列的表中,再批量更新该列的值。

步骤说明

  1. 确保目标SQL表已经存在file_name列(建议设为VARCHAR类型)
  2. 按原流程导入CSV数据(此时file_name列会为NULL)
  3. 导入完成后,执行UPDATE语句,将file_name列统一设置为当前去除后缀后的文件名

代码示例

// 原导入逻辑不变,导入完成后添加以下代码
const { Pool } = require('pg');

// 配置PostgreSQL连接(建议用Cloud SQL IAM认证或私有IP)
const pool = new Pool({
  user: process.env.DB_USER,
  host: process.env.DB_HOST,
  database: process.env.DB_NAME,
  password: process.env.DB_PASSWORD,
  port: 5432,
});

async function updateFileNameColumn(sqlTable, fileName) {
  const fileNameWithoutExt = path.parse(fileName).name;
  const query = `UPDATE ${sqlTable} SET file_name = $1 WHERE file_name IS NULL`;
  await pool.query(query, [fileNameWithoutExt]);
  await pool.end();
}

// 在原导入完成的Promise或回调后调用
// await updateFileNameColumn(sqlTable, csvFileName);

方案三:使用PostgreSQL COPY命令直接导入并添加列

如果你的Cloud SQL PostgreSQL实例支持访问Cloud Storage(需配置服务账号权限),可以直接用PostgreSQL的COPY命令结合SELECT语句,在导入时就添加文件名列。

步骤说明

  1. 给Cloud SQL的服务账号授予Cloud Storage的读取权限
  2. 使用gs://路径访问CSV文件,通过COPY ... FROM结合SELECT追加文件名列

代码示例

async function importWithFileName(bucketName, csvFileName, sqlTable) {
  const fileNameWithoutExt = path.parse(csvFileName).name;
  const csvPath = `gs://${bucketName}/${csvFileName}`;
  
  // 注意:需提前确保表结构包含file_name列
  const query = `
    COPY ${sqlTable} (col1, col2, file_name)
    FROM '${csvPath}'
    WITH (FORMAT csv, HEADER true)
    AS (col1 text, col2 text)
    SELECT col1, col2, '${fileNameWithoutExt}'::text AS file_name;
  `;

  await pool.query(query);
  await pool.end();
}

注意事项

  • 方案一需注意临时文件大小限制,Cloud Function默认临时磁盘空间为512MB,超大CSV建议用流式处理避免内存溢出。
  • 方案二要保证导入和更新的原子性,可包裹在事务中,避免中途失败导致数据不一致。
  • 方案三需要Cloud SQL PostgreSQL版本支持外部数据源,且需正确配置服务账号权限。

内容的提问来源于stack exchange,提问作者gboldt

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 22:20:34