【问题标题】:pg-promise: running a dependent query for each row in a query stream runs out of memorypg-promise:为查询流中的每一行运行依赖查询耗尽内存
【发布时间】:2021-02-04 19:10:40
【问题描述】:

在我的应用程序中,我需要为返回约 60k 行的查询中的每一行运行依赖更新。结果集太大而无法放入内存,因此自然的解决方案是流式传输结果并依次为每个结果运行相关查询。

无论我尝试什么,我的解决方案都内存不足,尽管我预计流式传输可以让我保持较低的内存使用率。

经过大量阅读和重读 SO、pg-promise Wiki 的各个页面,并尝试了不同的方法来实现这一点,我得出以下结论(简化我的代码):

try {
    const startTime = new Date();
    await db.tx("test-tx", async tx => {
        const qs = new QueryStream(`SELECT s.a AS i FROM GENERATE_SERIES(1, 100000) AS s(a)`);
        const result = await tx.stream(qs, stream => {
            return pgp.spex.stream.read(
                stream,
                async (i, row) => {
                    // console.log(`handling ${i}: ${JSON.stringify(data)}`);
                    await innerQuery(tx, row.i, startTime);
                },
                { readChunks: true }
            )
            .then(r => console.log("read done", r));
        });
        console.log("stream done", result);
    });
    console.log(`transaction done: ${memUsage()}MB, ${duration(startTime)} seconds`);
} catch (error) {
    console.error(error);
} finally {
    db.client.$pool.end();
}

async function innerQuery(tx, count, startTime) {
    if (count % 10000 === 0) {
        console.log(`row ${count}: ${memUsage()}MB, ${duration(startTime)} seconds`);
    }
    await tx.one("SELECT 1");
    if (count % 10000 === 0) {
        console.log(`inner query ${count} done`);
    }
}

function duration(startTime) {
    return Math.round((new Date() - startTime) / 1000);
}

function memUsage() {
    return Math.round(process.memoryUsage().heapUsed / 1024 / 1024);
}

这会按预期运行查询,但内存使用量一直在上升:

row 10000: 66MB, 1 seconds
row 20000: 124MB, 2 seconds
row 30000: 181MB, 3 seconds
row 40000: 241MB, 4 seconds
row 50000: 298MB, 5 seconds
row 60000: 355MB, 6 seconds
row 70000: 415MB, 7 seconds
row 80000: 474MB, 8 seconds
row 90000: 532MB, 9 seconds
row 100000: 593MB, 10 seconds
read done { calls: 100000, reads: 100000, length: 100000, duration: 10054 }
stream done { processed: 100000, duration: 10054 }
inner query 10000 done
inner query 20000 done
inner query 30000 done
inner query 40000 done
inner query 50000 done
inner query 60000 done
inner query 70000 done
inner query 80000 done
inner query 90000 done
inner query 100000 done
transaction done: 641MB, 42 seconds

现在,这里有一些东西:请注意,tx.stream 调用在所有内部查询解析之前返回,并打印到控制台。这解释了内存问题,所有这些闭包和承诺(其中 10 万个)都以某种方式在内存中等待流完成,以便它们自己解决并被 GC。

另一个数据点:如果我在顶层从db.tx 更改为db.task,在连接关闭之前只有一个或两个内部查询运行,进一步查询会导致错误(Querying against a released or lost connection.)。

我也尝试过使用tx.batch 和使用readChunks: false 进行stream.read 调用,但这只是在单个批处理后停止并锁定。

那我做错了什么?如何让内部查询在完成后立即解决,以便 GC 逐步回收内存?

【问题讨论】:

  • 我发现至少有一个问题 - return pgp.spex.stream.read。您在一个不期望返回任何承诺的回调中 - check API。这只是流初始化回调,但您正在从它返回一个承诺,然后丢失。这可能是问题的一部分。问题肯定是在某个地方产生了太多处于等待状态的 Promise。
  • 另外值得注意的是,在测试这么多记录时,使用console 是一个坏主意,因为控制台非常慢并且本身会消耗大量内存。除此之外,需要进行一些调试才能查看哪些承诺被泄露(如果有的话)。此外,也许readSize 可能会有一些用处。
  • 啊哈!我想我在关注这个例子:github.com/vitaly-t/pg-promise/wiki/Learn-by-Example#mixing-up 当一个承诺从stream 返回时。我会看看这是否对我有帮助。干杯

标签: node.js postgresql pg-promise node-streams


【解决方案1】:

据我所知,没有明显的方法可以减慢查询流以等待某些相关查询完成。新内部查询的创建速度与结果流式传输的速度一样快,因此内存消耗达到了顶峰。

我找到了一个不使用 QueryStream 的解决方案。这使用服务器端游标,意味着所有查询都按顺序运行。尚未探索尝试并行运行这些块以增加吞吐量,但它确实解决了内存问题。

const startTime = new Date();
await db.tx("test-tx", async tx => {
    await tx.none(`
        DECLARE test_cursor CURSOR FOR
        SELECT s.a AS i FROM GENERATE_SERIES(1, 100000) AS s(a)`);
    let row;
    while ((row = await tx.oneOrNone("FETCH NEXT FROM test_cursor"))) {
        await innerQuery(tx, row.i, startTime);
    }
    await tx.none("CLOSE test_cursor");
    console.log("outer query done");
});
console.log(`transaction done: ${memUsage()}MB, ${duration(startTime)} seconds`);

这会输出类似

row 10000: 9MB, 5 seconds
inner query 10000 done
row 20000: 8MB, 9 seconds
inner query 20000 done
row 30000: 10MB, 14 seconds
inner query 30000 done
row 40000: 9MB, 19 seconds
inner query 40000 done
row 50000: 11MB, 23 seconds
inner query 50000 done
row 60000: 10MB, 28 seconds
inner query 60000 done
row 70000: 8MB, 33 seconds
inner query 70000 done
row 80000: 11MB, 38 seconds
inner query 80000 done
row 90000: 9MB, 43 seconds
inner query 90000 done
row 100000: 12MB, 48 seconds
inner query 100000 done
outer query done
transaction done: 12MB, 48 seconds

【讨论】:

  • 是的,正确地流式传输数据一直是最棘手的部分。确保任何错误和内存泄漏。调试数据流可能很复杂。幸好您找到了适合您的案例的强大解决方案。
猜你喜欢
  • 1970-01-01
  • 2023-03-19
  • 1970-01-01
  • 2018-02-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多