三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

Node.js核心技术解析与实战考核指南

Node.js核心技术解析与实战考核指南

1. Node.js 考核题解析与实战指南

作为现代Web开发的核心技术之一,Node.js已经成为全栈工程师的必备技能。最近在技术社区看到一个很有意思的Node.js考核题集,这些题目不仅覆盖了基础知识点,更包含了实际开发中常见的场景化问题。今天我就结合自己多年的Node.js开发经验,带大家深度解析这些考核题背后的技术要点,并分享一些教科书上不会写的实战技巧。

2. Node.js基础概念考核

2.1 事件循环机制详解

Node.js最核心的特性就是其基于事件循环的非阻塞I/O模型。很多初学者对"非阻塞"的理解停留在表面,其实这涉及到三个关键组件:

  1. 调用栈(Call Stack):同步代码的执行场所
  2. 任务队列(Task Queue):存放异步操作完成后的回调函数
  3. 事件循环(Event Loop):不断检查调用栈和任务队列的协调者
// 典型的事件循环示例 console.log('Start'); setTimeout(() => { console.log('Timeout callback'); }, 0); Promise.resolve().then(() => { console.log('Promise resolved'); }); console.log('End');

这段代码的输出顺序是:

  1. Start
  2. End
  3. Promise resolved
  4. Timeout callback

关键点:微任务(Promise)优先级高于宏任务(setTimeout)

2.2 模块系统常见考点

CommonJS模块系统是Node.js的基础,考核中常出现以下题型:

  1. 模块缓存机制:同一个模块多次require只会执行一次
  2. 循环引用处理:Node.js通过未完成的exports对象解决
  3. 模块查找规则:node_modules逐级向上查找
// moduleA.js exports.value = 'A'; const b = require('./moduleB'); console.log('A中获取B的值:', b.value); exports.value = 'AA'; // moduleB.js exports.value = 'B'; const a = require('./moduleA'); console.log('B中获取A的值:', a.value); exports.value = 'BB';

这个循环引用的输出结果展示了Node.js模块加载的巧妙设计。

3. 异步编程考核实战

3.1 回调地狱与解决方案

早期的Node.js代码常陷入"回调地狱",现代开发中我们有多种解决方案:

  1. Promise链式调用
  2. async/await语法糖
  3. 第三方库如Bluebird
// 回调地狱示例 fs.readFile('file1.txt', (err, data1) => { if (err) throw err; fs.readFile('file2.txt', (err, data2) => { if (err) throw err; // 更多嵌套... }); }); // 使用async/await改进 async function readFiles() { try { const data1 = await fs.promises.readFile('file1.txt'); const data2 = await fs.promises.readFile('file2.txt'); // 更清晰的逻辑 } catch (err) { console.error(err); } }

3.2 错误处理最佳实践

Node.js中的错误处理有几个关键原则:

  1. 同步代码使用try/catch
  2. 异步回调遵循error-first约定
  3. Promise使用catch处理
  4. 全局错误监听process.on('uncaughtException')
// 错误处理综合示例 app.get('/api', async (req, res) => { try { const data = await someAsyncOperation(); res.json(data); } catch (err) { console.error('API错误:', err); res.status(500).send('服务器错误'); } }); // 全局错误捕获 process.on('uncaughtException', (err) => { console.error('未捕获异常:', err); // 优雅关闭进程 process.exit(1); });

4. 性能优化考核要点

4.1 内存泄漏排查

Node.js应用常见的内存泄漏场景:

  1. 全局变量滥用
  2. 闭包未释放
  3. 定时器未清除
  4. 事件监听未移除

使用以下工具进行诊断:

  • process.memoryUsage()
  • node --inspect+ Chrome DevTools
  • heapdump模块生成内存快照
// 典型内存泄漏示例 const leaks = []; app.get('/leak', (req, res) => { leaks.push(new Array(1000000).join('*')); res.send('内存增长中...'); });

4.2 集群模式优化

利用多核CPU的集群模式是Node.js性能优化的关键:

const cluster = require('cluster'); const os = require('os'); if (cluster.isMaster) { const cpuCount = os.cpus().length; for (let i = 0; i < cpuCount; i++) { cluster.fork(); } cluster.on('exit', (worker) => { console.log(`Worker ${worker.id} died`); cluster.fork(); }); } else { require('./app'); }

注意:共享状态需要使用Redis等外部存储

5. 实战项目考核解析

5.1 RESTful API设计

一个完整的Node.js API项目考核通常包含:

  1. 路由设计:遵循RESTful规范
  2. 中间件使用:身份验证、日志记录等
  3. 数据库集成:MongoDB/MySQL等
  4. 测试覆盖:单元测试和集成测试
// 典型API项目结构 const express = require('express'); const app = express(); // 中间件 app.use(express.json()); app.use(authMiddleware); // 路由 app.use('/api/users', require('./routes/users')); app.use('/api/products', require('./routes/products')); // 错误处理 app.use((err, req, res, next) => { res.status(500).json({ error: err.message }); }); // 启动 app.listen(3000, () => { console.log('Server running on port 3000'); });

5.2 WebSocket实时应用

现代Node.js考核常包含实时功能实现:

const WebSocket = require('ws'); const wss = new WebSocket.Server({ port: 8080 }); wss.on('connection', (ws) => { console.log('新客户端连接'); ws.on('message', (message) => { console.log('收到消息:', message); // 广播给所有客户端 wss.clients.forEach((client) => { if (client.readyState === WebSocket.OPEN) { client.send(message); } }); }); ws.on('close', () => { console.log('客户端断开连接'); }); });

6. 调试与测试技巧

6.1 高级调试技术

除了console.log,Node.js开发者应该掌握:

  1. 内置调试器node inspect
  2. Chrome DevTools集成
  3. VS Code调试配置
  4. 日志分级管理:winston/morgan等
// .vscode/launch.json配置示例 { "version": "0.2.0", "configurations": [ { "type": "node", "request": "launch", "name": "启动程序", "skipFiles": ["<node_internals>/**"], "program": "${workspaceFolder}/app.js" } ] }

6.2 测试金字塔实践

完整的测试策略应包含:

  1. 单元测试:Jest/Mocha
  2. 集成测试:Supertest
  3. E2E测试:Puppeteer
// Jest测试示例 const { sum } = require('./math'); describe('数学工具', () => { test('两数相加', () => { expect(sum(1, 2)).toBe(3); expect(sum(-1, 1)).toBe(0); }); test('非数字输入', () => { expect(() => sum('a', 1)).toThrow('参数必须是数字'); }); });

7. 安全防护要点

7.1 常见安全威胁防护

Node.js应用必须防范的安全风险:

  1. 注入攻击:SQL/NoSQL注入
  2. XSS攻击:输入输出过滤
  3. CSRF防护:CSRF令牌
  4. 敏感信息泄露:环境变量管理
// Helmet中间件提供基础安全防护 const helmet = require('helmet'); app.use(helmet()); // 防止NoSQL注入 app.use((req, res, next) => { const { query, body } = req; // 清理查询参数 const clean = (obj) => { Object.keys(obj).forEach(key => { if (typeof obj[key] === 'string') { obj[key] = obj[key].replace(/\$/g, ''); } }); }; clean(query); clean(body); next(); });

7.2 认证授权实现

现代认证方案实现:

  1. JWT实现:jsonwebtoken
  2. OAuth集成:passport.js
  3. 会话管理:express-session
// JWT认证示例 const jwt = require('jsonwebtoken'); function generateToken(user) { return jwt.sign( { userId: user.id }, process.env.JWT_SECRET, { expiresIn: '1h' } ); } function authenticate(req, res, next) { const token = req.headers.authorization?.split(' ')[1]; if (!token) return res.status(401).send('需要认证'); try { const decoded = jwt.verify(token, process.env.JWT_SECRET); req.userId = decoded.userId; next(); } catch (err) { res.status(401).send('无效令牌'); } }

8. 部署与监控

8.1 生产环境部署

Node.js应用部署要点:

  1. 进程管理:PM2/Nodemon
  2. 反向代理:Nginx配置
  3. 日志管理:ELK Stack
  4. 容器化:Docker部署
# PM2常用命令 pm2 start app.js -i max # 集群模式启动 pm2 logs # 查看日志 pm2 monit # 监控面板 pm2 save # 保存进程列表 pm2 startup # 设置开机启动

8.2 性能监控方案

生产环境监控策略:

  1. 指标收集:Prometheus
  2. 可视化:Grafana
  3. 告警系统:Alertmanager
  4. APM工具:New Relic/AppDynamics
// 使用prom-client收集指标 const client = require('prom-client'); const collectDefaultMetrics = client.collectDefaultMetrics; collectDefaultMetrics({ timeout: 5000 }); app.get('/metrics', async (req, res) => { res.set('Content-Type', client.register.contentType); res.end(await client.register.metrics()); });

9. 最新特性与趋势

9.1 ES模块支持

Node.js对ES模块的原生支持:

// package.json { "type": "module" } // app.js import express from 'express'; import { readFile } from 'fs/promises'; const app = express(); const data = await readFile('config.json');

9.2 Worker Threads应用

CPU密集型任务的解决方案:

const { Worker, isMainThread } = require('worker_threads'); if (isMainThread) { // 主线程 const worker = new Worker(__filename); worker.on('message', (msg) => { console.log('来自Worker的消息:', msg); }); worker.postMessage('主线程消息'); } else { // Worker线程 parentPort.on('message', (msg) => { console.log('来自主线程的消息:', msg); parentPort.postMessage('Worker回复'); }); }

10. 综合实战案例分析

10.1 文件上传服务

完整文件上传实现:

const multer = require('multer'); const path = require('path'); // 存储配置 const storage = multer.diskStorage({ destination: (req, file, cb) => { cb(null, 'uploads/'); }, filename: (req, file, cb) => { const ext = path.extname(file.originalname); cb(null, `${Date.now()}${ext}`); } }); // 文件过滤 const fileFilter = (req, file, cb) => { const allowedTypes = ['image/jpeg', 'image/png']; if (allowedTypes.includes(file.mimetype)) { cb(null, true); } else { cb(new Error('不支持的文件类型'), false); } }; const upload = multer({ storage, fileFilter }); app.post('/upload', upload.single('file'), (req, res) => { if (!req.file) { return res.status(400).send('请上传有效文件'); } res.json({ url: `/uploads/${req.file.filename}` }); });

10.2 定时任务系统

基于Node.js的定时任务实现:

const schedule = require('node-schedule'); const axios = require('axios'); // 每小时执行 schedule.scheduleJob('0 * * * *', async () => { console.log('开始执行定时任务'); try { const res = await axios.get('https://api.example.com/data'); // 处理数据... } catch (err) { console.error('定时任务失败:', err); } }); // 复杂规则 const rule = new schedule.RecurrenceRule(); rule.dayOfWeek = [0, new schedule.Range(1, 5)]; // 周一到周五 rule.hour = 9; rule.minute = 30; schedule.scheduleJob(rule, () => { console.log('工作日早上9:30执行'); });

11. 性能调优实战

11.1 数据库查询优化

Node.js + MongoDB性能优化技巧:

// 低效查询 const users = await User.find({}) .skip((page - 1) * limit) .limit(limit) .sort({ createdAt: -1 }); // 优化后 const users = await User.find({}, 'name email avatar') // 只选择必要字段 .lean() // 返回普通JS对象而非Mongoose文档 .skip((page - 1) * limit) .limit(limit) .sort({ createdAt: -1 }) .cache({ key: `users_page_${page}` }); // 添加缓存

11.2 流处理大数据

高效处理大文件的流式方法:

const fs = require('fs'); const zlib = require('zlib'); // 传统方式 - 内存占用高 fs.readFile('large.log', (err, data) => { if (err) throw err; const compressed = zlib.gzipSync(data); fs.writeFile('large.log.gz', compressed, (err) => { if (err) throw err; }); }); // 流式处理 - 内存高效 fs.createReadStream('large.log') .pipe(zlib.createGzip()) .pipe(fs.createWriteStream('large.log.gz')) .on('finish', () => { console.log('压缩完成'); });

12. 微服务架构实现

12.1 gRPC服务通信

Node.js实现gRPC微服务:

// service.proto syntax = "proto3"; service ProductService { rpc GetProduct (ProductRequest) returns (ProductResponse); } message ProductRequest { string id = 1; } message ProductResponse { string id = 1; string name = 2; float price = 3; }
// 服务端实现 const grpc = require('@grpc/grpc-js'); const protoLoader = require('@grpc/proto-loader'); const packageDefinition = protoLoader.loadSync('service.proto'); const proto = grpc.loadPackageDefinition(packageDefinition); const server = new grpc.Server(); server.addService(proto.ProductService.service, { GetProduct: (call, callback) => { const product = getProductFromDB(call.request.id); callback(null, product); } }); server.bindAsync( '0.0.0.0:50051', grpc.ServerCredentials.createInsecure(), (err, port) => { server.start(); } );

12.2 服务健康检查

生产级健康检查实现:

const health = require('express-healthcheck'); app.use('/health', health({ healthy: () => { return { status: 'up', db: checkDatabaseConnection(), redis: checkRedisConnection(), uptime: process.uptime() }; } })); // 自定义检查函数 async function checkDatabaseConnection() { try { await mongoose.connection.db.admin().ping(); return 'connected'; } catch (err) { return 'disconnected'; } }

13. 错误监控与日志

13.1 结构化日志实践

生产环境日志最佳实践:

const winston = require('winston'); const { ElasticsearchTransport } = require('winston-elasticsearch'); const logger = winston.createLogger({ level: 'info', format: winston.format.combine( winston.format.timestamp(), winston.format.json() ), transports: [ new winston.transports.Console(), new winston.transports.File({ filename: 'combined.log' }), new ElasticsearchTransport({ level: 'info', clientOpts: { node: 'http://localhost:9200' } }) ] }); // 使用示例 logger.info('用户登录', { userId: 123, ip: '192.168.1.1' }); logger.error('数据库连接失败', { error: err.stack });

13.2 分布式追踪实现

使用OpenTelemetry实现分布式追踪:

const { NodeTracerProvider } = require('@opentelemetry/node'); const { SimpleSpanProcessor } = require('@opentelemetry/tracing'); const { JaegerExporter } = require('@opentelemetry/exporter-jaeger'); const provider = new NodeTracerProvider(); provider.register(); const exporter = new JaegerExporter({ serviceName: 'nodejs-service', host: 'jaeger-agent' }); provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); // 自动instrumentation常用库 require('@opentelemetry/plugin-http'); require('@opentelemetry/plugin-express'); require('@opentelemetry/plugin-mongodb');

14. Serverless架构应用

14.1 AWS Lambda部署

Node.js Lambda函数最佳实践:

// handler.js exports.handler = async (event, context) => { try { const data = JSON.parse(event.body); // 业务逻辑处理 const result = await processData(data); return { statusCode: 200, body: JSON.stringify(result) }; } catch (err) { return { statusCode: 500, body: JSON.stringify({ error: err.message }) }; } }; // 本地测试 if (process.env.NODE_ENV === 'development') { const event = { body: JSON.stringify({ test: 'data' }) }; exports.handler(event) .then(console.log) .catch(console.error); }

14.2 冷启动优化

Serverless性能优化技巧:

  1. 减小包体积:只包含必要依赖
  2. 预初始化:在handler外初始化资源
  3. 保持活跃:定时ping函数
  4. 使用Provisioned Concurrency
// 预初始化数据库连接 const mongoose = require('mongoose'); let conn = null; const uri = process.env.MONGODB_URI; exports.handler = async (event, context) => { context.callbackWaitsForEmptyEventLoop = false; if (!conn) { conn = await mongoose.createConnection(uri, { bufferCommands: false, bufferMaxEntries: 0, useNewUrlParser: true, useUnifiedTopology: true }); } const Model = conn.model('Test', new mongoose.Schema({ name: String })); const docs = await Model.find(); return { statusCode: 200, body: JSON.stringify(docs) }; };

15. 现代全栈开发模式

15.1 GraphQL API实现

Node.js + GraphQL全栈方案:

const { ApolloServer, gql } = require('apollo-server-express'); const typeDefs = gql` type Query { users: [User!]! user(id: ID!): User } type User { id: ID! name: String! email: String! posts: [Post!]! } type Post { id: ID! title: String! content: String! } `; const resolvers = { Query: { users: () => User.find(), user: (_, { id }) => User.findById(id) }, User: { posts: (parent) => Post.find({ userId: parent.id }) } }; const server = new ApolloServer({ typeDefs, resolvers }); server.applyMiddleware({ app });

15.2 同构渲染应用

Next.js服务端渲染集成:

// pages/api/users.js export default async function handler(req, res) { const response = await fetch('http://localhost:3000/api/graphql', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ query: ` query { users { id name } } ` }) }); const { data } = await response.json(); res.status(200).json(data.users); } // pages/users.js export async function getServerSideProps() { const res = await fetch('http://localhost:3000/api/users'); const users = await res.json(); return { props: { users } }; } export default function UsersPage({ users }) { return ( <ul> {users.map(user => ( <li key={user.id}>{user.name}</li> ))} </ul> ); }

16. 工具链与开发效率

16.1 现代化开发工具

提升Node.js开发效率的工具集:

  1. 调试工具:ndb、node-inspect
  2. 代码质量:ESLint + Prettier
  3. 类型检查:TypeScript/JSDoc
  4. API测试:Postman/Insomnia
  5. 性能分析:clinic.js
// 使用clinic.js进行性能分析 const clinic = require('clinic'); async function runProfiler() { const doctor = new clinic.Doctor(); await doctor.collect(['node', 'server.js'], { detectPort: true, sampleInterval: 100 }); await doctor.visualize(); } runProfiler().catch(console.error);

16.2 自动化部署流水线

CI/CD最佳实践配置:

# .github/workflows/deploy.yml name: Node.js CI/CD on: push: branches: [ main ] pull_request: branches: [ main ] jobs: test: runs-on: ubuntu-latest steps: - uses: actions/checkout@v2 - name: Use Node.js uses: actions/setup-node@v2 with: node-version: '16' - run: npm ci - run: npm test deploy: needs: test runs-on: ubuntu-latest steps: - uses: actions/checkout@v2 - name: Use Node.js uses: actions/setup-node@v2 with: node-version: '16' - run: npm ci - run: npm run build - name: Deploy to AWS uses: aws-actions/configure-aws-credentials@v1 with: aws-access-key-id: ${{ secrets.AWS_ACCESS_KEY_ID }} aws-secret-access-key: ${{ secrets.AWS_SECRET_ACCESS_KEY }} aws-region: us-east-1 - name: Install Serverless run: npm install -g serverless - name: Deploy run: sls deploy --stage production

17. 性能基准测试

17.1 负载测试实践

使用Artillery进行压力测试:

# load-test.yml config: target: "http://localhost:3000" phases: - duration: 60 arrivalRate: 10 name: "Warm up" - duration: 120 arrivalRate: 20 rampTo: 100 name: "Ramp up load" - duration: 300 arrivalRate: 100 name: "Sustained load" scenarios: - name: "API测试" flow: - get: url: "/api/products" - post: url: "/api/orders" json: productId: "123" quantity: 2

17.2 性能优化指标

关键性能指标(KPI)监控:

  1. 吞吐量:RPS(Requests Per Second)
  2. 延迟:P95/P99响应时间
  3. 错误率:HTTP 5xx比例
  4. 资源使用:CPU/内存占用
// 性能监控中间件 function performanceMiddleware(req, res, next) { const start = process.hrtime(); res.on('finish', () => { const duration = process.hrtime(start); const ms = duration[0] * 1000 + duration[1] / 1e6; metrics.httpRequestsTotal.inc(); metrics.httpRequestDurationMicroseconds.observe(ms); if (res.statusCode >= 500) { metrics.httpErrorsTotal.inc(); } }); next(); }

18. 高级设计模式

18.1 依赖注入实现

Node.js中的IoC容器实践:

// container.js class Container { constructor() { this.services = {}; } register(name, callback) { this.services[name] = callback; } resolve(name) { if (!this.services[name]) { throw new Error(`服务未注册: ${name}`); } if (typeof this.services[name] === 'function') { this.services[name] = this.services[name](this); } return this.services[name]; } } // 使用示例 const container = new Container(); container.register('db', () => { return new Database(process.env.DB_URL); }); container.register('userService', (c) => { return new UserService(c.resolve('db')); }); const userService = container.resolve('userService');

18.2 CQRS模式应用

命令查询职责分离实现:

// commands/AddProductCommand.js class AddProductCommand { constructor(productData) { this.productData = productData; } async execute() { const product = new Product(this.productData); await product.save(); await eventBus.publish('ProductAdded', product); return product; } } // queries/GetProductsQuery.js class GetProductsQuery { constructor(filters) { this.filters = filters; } async execute() { return Product.find(this.filters).lean(); } } // 使用示例 router.post('/products', async (req, res) => { const command = new AddProductCommand(req.body); const product = await command.execute(); res.status(201).json(product); }); router.get('/products', async (req, res) => { const query = new GetProductsQuery(req.query); const products = await query.execute(); res.json(products); });

19. 微前端集成方案

19.1 服务端集成微前端

Node.js作为微前端聚合层:

const express = require('express'); const { createProxyMiddleware } = require('http-proxy-middleware'); const app = express(); // 静态资源服务 app.use('/static', express.static('public')); // 微前端路由 app.use('/app1', createProxyMiddleware({ target: 'http://app1-service', pathRewrite: { '^/app1': '' }, changeOrigin: true })); app.use('/app2', createProxyMiddleware({ target: 'http://app2-service', pathRewrite: { '^/app2': '' }, changeOrigin: true })); // 聚合页面 app.get('*', async (req, res) => { const [app1Header, app2Content] = await Promise.all([ fetch('http://app1-service/header'), fetch('http://app2-service/content') ]); res.send(` <html> <head> <title>微前端聚合</title> </head> <body> <div id="header">${await app1Header.text()}</div> <div id="content">${await app2Content.text()}</div> </body> </html> `); });

19.2 模块联邦应用

Webpack Module Federation实现:

// app1/webpack.config.js const ModuleFederationPlugin = require('webpack/lib/container/ModuleFederationPlugin'); module.exports = { plugins: [ new ModuleFederationPlugin({ name: 'app1', filename: 'remoteEntry.js', exposes: { './Header': './src/Header' }, shared: ['react', 'react-dom'] }) ] }; // app2/webpack.config.js const ModuleFederationPlugin = require('webpack/lib/container/ModuleFederationPlugin'); module.exports = { plugins: [ new ModuleFederationPlugin({ name: 'app2', remotes: { app1: 'app1@http://localhost:3001/remoteEntry.js' }, shared: ['react', 'react-dom'] }) ] }; // app2/src/App.js import React from 'react'; const Header = React.lazy(() => import('app1/Header')); function App() { return ( <React.Suspense fallback="Loading..."> <Header /> <div>App2 Content</div> </React.Suspense> ); }

20. 边缘计算应用

20.1 Edge Functions实现

基于Node.js的边缘函数:

// 边缘路由逻辑 addEventListener('fetch', event => { event.respondWith(handleRequest(event.request)); }); async function handleRequest(request) { const url = new URL(request.url); // A/B测试路由 if (url.pathname.startsWith('/product')) { const variant = Math.random() > 0.5 ? 'a' : 'b'; return fetch(`https://${variant}.cdn.example.com${url.pathname}`); } // 地理路由 const country = request.headers.get('cf-ipcountry'); if (country === 'CN') { return Response.redirect('https://cn.example.com', 302); } // 默认回源 return fetch(`https://origin.example.com${url.pathname}`); }

20.2 边缘缓存策略

CDN边缘缓存优化:

// 缓存控制中间件 function edgeCacheControl(req, res, next) { // 静态资源缓存1年 if (req.path.startsWith('/static')) { res.set('Cache-Control', 'public, max-age=31536000, immutable'); res.set('CDN-Cache-Control', 'public, s-maxage=31536000'); } // API响应缓存5分钟 else if (req.path.startsWith('/api')) { res.set('Cache-Control', 'no-cache'); res.set('CDN-Cache-Control', 'public, s-maxage=300'); } // 页面级缓存1小时 else { res.set('Cache-Control', 'no-cache'); res.set('CDN-Cache-Control', 'public, s-maxage=3600'); } next(); }

21. 实时数据处理

21.1 WebSocket集群

分布式WebSocket服务:

const WebSocket = require('ws'); const Redis = require('ioredis'); const pub = new Redis(process.env.REDIS_URL); const sub = new Redis(process.env.REDIS_URL); const wss = new WebSocket.Server({ port: 8080 }); // 订阅Redis频道 sub.subscribe('messages'); wss.on('connection', (ws) => { console.log('新客户端连接'); // 广播消息到Redis ws.on('message', (message) => { pub.publish('messages', message); }); // 从Redis接收广播 sub.on('message', (channel, message) => { if (channel === 'messages') { ws.send(message); } }); ws.on('close', () => { console.log('客户端断开连接'); }); });

21.2 实时数据分析

使用Node.js处理实时数据流:

const { Kafka } = require('kafkajs'); const kafka = new Kafka({ clientId: 'nodejs-consumer', brokers: ['kafka1:9092', 'kafka2:9092'] }); const consumer = kafka.consumer({ groupId: 'analytics-group' }); async function run() { await consumer.connect(); await consumer.subscribe({ topic: 'user-events' }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { const event = JSON.parse(message.value.toString()); // 实时分析逻辑 if (event.type === 'page_view') { await updatePageStats(event.url); } else if (event.type === 'purchase') { await updateRevenueStats(event.amount); } } }); } run().catch(console.error);

22. AI集成应用

22.1 TensorFlow.js集成

Node.js中的机器学习:

const tf = require('@tensorflow/tfjs-node'); const fs = require('fs').promises; async function loadModel() { const model = await tf.loadLayersModel('file://./model/model.json'); return model; } async function predict(imagePath) { const model = await loadModel(); // 读取并预处理图像 const imageBuffer = await fs.readFile(imagePath); const tfimage = tf.node.decodeImage(imageBuffer); const input = tfimage.resizeBilinear([224, 224]) .div(255.0) .expandDims(); // 预测 const predictions = model.predict(input); const results = await predictions.array(); return results[0]; } // 使用示例 predict('test.jpg').then(results => { console.log('预测结果:', results); });

22.2 自然语言处理

使用NLP库处理文本:

const natural
← 返回列表