Web3监听合约事件写数据库并发时数据串扰问题
问题原因
- 核心bug是所有事件关联的业务字段(
name/symbol/collectionAddress/addressWallet/walletID/tokenURI/description)都定义在全局作用域,所有事件回调共享这一块内存。多用户并发触发事件时,后执行的工厂事件回调会直接覆盖前一个事件流程还没完成入库的变量值,等对应NFT合约的铸造事件触发、执行数据库写入时,读到的已经是其他用户的事件数据,自然出现串扰错乱。 - 额外存在两个隐患:
- 每次收到工厂的新合约创建事件就给对应NFT合约绑定一次事件监听,没有做去重,Web3 WS重连拉取历史事件时会重复绑定回调,导致重复入库。
- 数据库操作用字符串拼接SQL,既存在注入风险,也因为回调式写法没有做异步流程控制,并发场景下数据库操作执行顺序不可控。
修复方案
- 移除所有全局业务变量,单次事件流程的所有数据通过参数、闭包传递,保证每个事件触发时的上下文完全独立,不会互相覆盖。
- 增加已监听合约地址缓存,避免重复给同一个NFT合约绑定事件监听。
- 把数据库操作封装为Promise形式,改用参数化查询,保证异步流程执行顺序可控。
修复后的完整代码如下:
const db = require("../connectDB/DB"); const Web3 = require("web3"); const HDWalletProvider = require("@truffle/hdwallet-provider"); const EzeyNFTABI = require("../../client/src/Artifacts/EzeyNFT.json"); const EzeyNFTFactoryABI = require("../../client/src/Artifacts/EzeyNFTFactory.json"); const wsProvider = new Web3.providers.WebsocketProvider( "wss://morning-twilight-cherry.matic-testnet.quiknode.pro/6ba9d2c5b8a046814b28f974c3643c679914f7ff/" ); HDWalletProvider.prototype.on = wsProvider.on.bind(wsProvider); let provider = new HDWalletProvider( "PRIVAE-KEY", wsProvider ); const web3 = new Web3(provider); // 仅做监听地址缓存,不存业务数据 const listenedCollection = new Set(); // 封装db查询为Promise function dbQuery(sql, params = []) { return new Promise((resolve, reject) => { db.query(sql, params, (err, result) => { if (err) return reject(err); resolve(result); }) }) } // 入库逻辑接收完整上下文参数,不读全局变量 async function insertToCollectionTable(collectionData) { const { symbol, tokenURI, walletID, addressWallet } = collectionData; // 先查用户是否存在 const userResult = await dbQuery( `SELECT * FROM ezeyNFT.userNFT WHERE addressWallet = ?`, [addressWallet] ); // 用户不存在先创建 if (userResult.length === 0) { await dbQuery( `INSERT INTO ezeyNFT.userNFT (id,addressWallet) VALUES (?, ?)`, [walletID, addressWallet] ); } // 写入NFT集合数据 await dbQuery( `INSERT INTO ezeyNFT.collectionNFT (NFTSymbol,NFTUrl,IDAddressWallet) VALUES (?,?,?)`, [symbol, tokenURI, walletID] ); } // mint事件回调,透传工厂事件拿到的上下文 function getMintEventHandler(collectionContext) { return async function (err, data) { if (err) { console.log(err); return; } const result = data.returnValues; const fullData = { ...collectionContext, tokenURI: result.tokenURI, description: result.description } console.log('mint event:', result); await insertToCollectionTable(fullData); } } async function bindCollectionListener(collectionContext) { const { collectionAddress } = collectionContext; // 已经监听过的地址直接跳过,避免重复绑定 if (listenedCollection.has(collectionAddress)) return; listenedCollection.add(collectionAddress); const nftContract = new web3.eth.Contract( EzeyNFTABI.abi, collectionAddress ); nftContract.events.newCollection(getMintEventHandler(collectionContext)); } async function eventFactory_Handler(err, data) { if (err) { console.log(err); return; } const result = data.returnValues; // 本次事件的上下文存在函数作用域,不赋值给全局变量 const collectionContext = { name: result.name, symbol: result.symbol, collectionAddress: result.collectionAddress, addressWallet: result.addressWallet, walletID: result.walletID } console.log('factory event:', result); await bindCollectionListener(collectionContext); } async function startListener() { const factoryContract = new web3.eth.Contract( EzeyNFTFactoryABI.abi, "0x0a77174a8F78E64fFE589288D2C47AE83189dEAd" ); factoryContract.events.newCollection(eventFactory_Handler); } startListener(); module.exports = { insertToCollectionTable, };
注意:如果需要处理WS断连重连场景,可以给
wsProvider加end/error事件监听,重连后清空listenedCollection缓存重新绑定监听即可,避免历史事件重复触发。
内容的提问来源于stack exchange,提问作者asalef alena
相关产品推荐
相关产品推荐

