RabbitMQ HeartBeat Timeout issue

Viewed 627

I'm currently using RabbitMQ as a message broker. Recently, I see many error HeartBeat Timeout in my error log.

Also in RabbitMQ log, I see this log:

enter image description here

I don't know why there is too many connection from vary ranges of port. I use default setup without any further configuration. Here is my code used to publish and consume:

import { connect } from 'amqplib/callback_api';
import hanlder from '../calculator/middleware';
import { logger } from '../config/logger';

async function consumeRabbitMQServer(serverURL, exchange, queue) {
  connect('amqp://localhost', async (error0, connection) => {
    if (error0) throw error0;

    const channel = connection.createChannel((error1) => {
      if (error1) throw error1;
    });

    channel.assertExchange(exchange, 'direct', {
      durable: true
    });

    channel.assertQueue(
      queue,
      {
        durable: true
      },
      (error2) => {
        if (error2) throw error2;
        logger.info(`Connect to ${serverURL} using queue ${queue}`);
      }
    );

    channel.prefetch(1);

    channel.bindQueue(queue, exchange, 'info');
    channel.noAck = true;

    channel.consume(queue, (msg) => {
      hanlder(JSON.parse(msg.content.toString()))
        .then(() => {
          channel.ack(msg);
        })
        .catch((err) => {
          channel.reject(msg);
        });
    });
  });
}

export default consumeRabbitMQServer;

Code used to publish message:

import createConnection from './connection';
import { logger } from '../config/logger';

async function publishToRabbitMQServer(serverURL, exchange, queue) {
  const connection = createConnection(serverURL);
  const c = await connection.then(async (conn) => {
    const channel = await conn.createChannel((error1) => {
      if (error1) throw error1;
    });

    channel.assertExchange(exchange, 'direct', {
      durable: true
    });

    channel.assertQueue(
      queue,
      {
        durable: true
      },
      (error2) => {
        if (error2) throw error2;
        logger.info(`Publish to ${serverURL} using queue ${queue}`);
      }
    );

    channel.bindQueue(queue, exchange, 'info');

    return channel;
  });
  return c;
}

export default publishToRabbitMQServer;

Whenever I start my server, I run this piece of code to create a client consume to RabbitMQ:

const { RABBITMQ_SERVER } = process.env;
consumeRabbitMQServer(RABBITMQ_SERVER, 'abc', 'abc');

And this piece of code is used when ever a message in need published to RabbitMQ

  const payloads = call.request.payloads;
  const { RABBITMQ_SERVER } = process.env;
  const channel = await publishToRabbitMQServer(RABBITMQ_SERVER, 'abc', 'abc');

  for (let i = 0; i < payloads.length; i++) {
    channel.publish('abc', 'info', Buffer.from(JSON.stringify(payloads[i])));
  }

I'm reusing code from RabbitMQ document, and it seem that this problem happen whenever there are too many user publish message. Thanks for helping. Update: I think the root cause is when I need to publish a message, I create a new connection. I'm working to improve it, any help is appreciate. Many thanks.

0 Answers
Related