如何同时配置NestJS应用,使其既支持Kafka微服务消费又能实现基于Postgres的CRUD REST接口
Got it, let's break this down step by step! You want your NestJS app to handle both Kafka message consumption and expose REST CRUD endpoints with Postgres integration—this is totally achievable. Here's exactly how to set it up:
1. Run NestJS as Both HTTP Server and Kafka Consumer
Your current main.ts only spins up a Kafka microservice, which is why you can't access REST endpoints. Instead, we'll create a regular HTTP app first, then attach the Kafka microservice to it.
Update your main.ts like this:
import { NestFactory } from '@nestjs/core'; import { MicroserviceOptions, Transport } from '@nestjs/microservices'; import { AppModule } from './app.module'; async function bootstrap() { // Create a standard HTTP app (defaults to port 3000) const app = await NestFactory.create(AppModule); // Attach your Kafka microservice configuration (keep your existing consumer setup) app.connectMicroservice<MicroserviceOptions>({ transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'], }, consumer: { groupId: 'groupId', }, }, }); // Start both the HTTP server and Kafka microservice await app.startAllMicroservices(); await app.listen(3000); // Adjust port if needed, e.g., 3001 console.log('✅ NestJS running at http://localhost:3000 and listening to Kafka'); } bootstrap();
Now your app will listen for REST requests on port 3000 and consume Kafka messages as before.
2. Set Up Postgres & TypeORM for CRUD
First, install the required dependencies for Postgres and TypeORM (NestJS's preferred ORM for database interactions):
npm install @nestjs/typeorm typeorm pg
Next, configure TypeORM in your AppModule to connect to your Postgres database:
import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { AppController } from './app.controller'; import { AppService } from './app.service'; // Import your Kafka consumer service here (e.g., KafkaConsumerService) import { UserModule } from './user/user.module'; // We'll create this CRUD module next @Module({ imports: [ // Configure Postgres connection TypeOrmModule.forRoot({ type: 'postgres', host: 'localhost', port: 5432, username: 'your-db-username', // Replace with your Postgres username password: 'your-db-password', // Replace with your Postgres password database: 'your-db-name', // Replace with your database name (create it first in Postgres) entities: [__dirname + '/**/*.entity{.ts,.js}'], // Auto-load all entity files synchronize: true, // Only use in development! Auto-creates tables (turn off in production) autoLoadEntities: true, }), UserModule, // Add your CRUD module here ], controllers: [AppController], providers: [AppService], // Add your Kafka consumer service to providers here }) export class AppModule {}
3. Build Your CRUD Modules/Controllers/Services
Let's create a sample User CRUD setup (you can adapt this to your own entity):
Step 3.1: Create the Entity (Database Model)
Create src/user/user.entity.ts:
import { Entity, Column, PrimaryGeneratedColumn } from 'typeorm'; @Entity() export class User { @PrimaryGeneratedColumn() id: number; @Column() name: string; @Column({ unique: true }) email: string; @Column() age: number; }
Step 3.2: Create the CRUD Module
Create src/user/user.module.ts:
import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { User } from './user.entity'; import { UserController } from './user.controller'; import { UserService } from './user.service'; @Module({ imports: [TypeOrmModule.forFeature([User])], // Register the User entity for this module controllers: [UserController], providers: [UserService], }) export class UserModule {}
Step 3.3: Create the CRUD Service
Create src/user/user.service.ts:
import { Injectable, NotFoundException } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { User } from './user.entity'; @Injectable() export class UserService { constructor( @InjectRepository(User) private userRepository: Repository<User>, ) {} // Create a new user async create(userData: Partial<User>): Promise<User> { const newUser = this.userRepository.create(userData); return await this.userRepository.save(newUser); } // Get all users async findAll(): Promise<User[]> { return await this.userRepository.find(); } // Get a single user by ID async findOne(id: number): Promise<User> { const user = await this.userRepository.findOneBy({ id }); if (!user) { throw new NotFoundException(`User with ID ${id} not found`); } return user; } // Update a user async update(id: number, userData: Partial<User>): Promise<User> { await this.userRepository.update(id, userData); return await this.findOne(id); } // Delete a user async remove(id: number): Promise<void> { const result = await this.userRepository.delete(id); if (result.affected === 0) { throw new NotFoundException(`User with ID ${id} not found`); } } }
Step 3.4: Create the REST Controller
Create src/user/user.controller.ts:
import { Controller, Get, Post, Put, Delete, Body, Param, HttpCode, HttpStatus } from '@nestjs/common'; import { UserService } from './user.service'; import { User } from './user.entity'; @Controller('users') // This sets the base route: http://localhost:3000/users export class UserController { constructor(private readonly userService: UserService) {} // Create user: POST /users @Post() create(@Body() userData: Partial<User>): Promise<User> { return this.userService.create(userData); } // Get all users: GET /users @Get() findAll(): Promise<User[]> { return this.userService.findAll(); } // Get single user: GET /users/:id @Get(':id') findOne(@Param('id') id: string): Promise<User> { return this.userService.findOne(+id); // Convert string ID to number } // Update user: PUT /users/:id @Put(':id') update(@Param('id') id: string, @Body() userData: Partial<User>): Promise<User> { return this.userService.update(+id, userData); } // Delete user: DELETE /users/:id @Delete(':id') @HttpCode(HttpStatus.NO_CONTENT) remove(@Param('id') id: string): Promise<void> { return this.userService.remove(+id); } }
4. Verify Kafka Consumer Still Works
Your existing Kafka consumer logic will continue to work as long as you keep it registered in the AppModule providers. For example, if you're using a dedicated consumer service:
import { Injectable, OnModuleInit } from '@nestjs/common'; import { KafkaService } from '@nestjs/microservices'; import { Consumer } from 'kafkajs'; @Injectable() export class KafkaConsumerService implements OnModuleInit { private consumer: Consumer; constructor(private readonly kafkaService: KafkaService) { this.consumer = this.kafkaService.getConsumer('groupId'); // Match your groupId from main.ts } async onModuleInit() { await this.consumer.subscribe({ topic: 'your-spring-boot-topic', fromBeginning: true }); await this.consumer.run({ eachMessage: async ({ topic, message }) => { console.log(`Received message from Kafka topic ${topic}:`, message.value.toString()); // Optional: Use your CRUD service to save the message to Postgres // await this.userService.create({...}); }, }); } }
Add this service to your AppModule providers array, and it will start consuming messages when the app boots.
Testing Everything
- Make sure Postgres is running, and you've created the database specified in
AppModule. - Start your NestJS app with
npm run start:dev. - Test the REST endpoints using tools like Postman or curl:
POST http://localhost:3000/userswith body{"name": "John Doe", "email": "john@example.com", "age": 30}GET http://localhost:3000/usersto see all users
- Send a message from your Spring Boot producer to your Kafka topic—check the NestJS console to confirm it's received.
内容的提问来源于stack exchange,提问作者prabath shalitha

