Browse Source

Merge pull request #13402 from abpframework/liangshiwei/rabbitmq

Enhance RabbitMQ to set PrefetchCount
pull/13405/head
maliming 4 years ago
committed by GitHub
parent
commit
1c74dbc09d
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 12
      docs/en/Background-Jobs-RabbitMq.md
  2. 3
      docs/en/Distributed-Event-Bus-RabbitMQ-Integration.md
  3. 10
      docs/zh-Hans/Background-Jobs-RabbitMq.md
  4. 3
      docs/zh-Hans/Distributed-Event-Bus-RabbitMQ-Integration.md
  5. 2
      framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/AbpRabbitMqBackgroundJobOptions.cs
  6. 10
      framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs
  7. 6
      framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueueConfiguration.cs
  8. 2
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/AbpRabbitMqEventBusOptions.cs
  9. 3
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  10. 6
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs
  11. 7
      framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs

12
docs/en/Background-Jobs-RabbitMq.md

@ -126,25 +126,31 @@ By default, all the job types use the `Default` RabbitMQ connection.
Configure<AbpRabbitMqBackgroundJobOptions>(options => Configure<AbpRabbitMqBackgroundJobOptions>(options =>
{ {
options.DefaultQueueNamePrefix = "my_app_jobs."; options.DefaultQueueNamePrefix = "my_app_jobs.";
options.DefaultDelayedQueueNamePrefix = "my_app_jobs.delayed"
options.PrefetchCount = 1;
options.JobQueues[typeof(EmailSendingArgs)] = options.JobQueues[typeof(EmailSendingArgs)] =
new JobQueueConfiguration( new JobQueueConfiguration(
typeof(EmailSendingArgs), typeof(EmailSendingArgs),
queueName: "my_app_jobs.emails", queueName: "my_app_jobs.emails",
connectionName: "SecondConnection" connectionName: "SecondConnection",
delayedQueueName:"my_app_jobs.emails.delayed"
); );
}); });
```` ````
* This example sets the default queue name prefix to `my_app_jobs.`. If different applications use the same RabbitMQ server, it would be important to use different prefixes for each application to not consume jobs of each other. * This example sets the default queue name prefix to `my_app_jobs.` and default delayed queue name prefix to `my_app_jobs.delayed`. If different applications use the same RabbitMQ server, it would be important to use different prefixes for each application to not consume jobs of each other.
* Sets `PrefetchCount` for all queues.
* Also specifies a different connection string for the `EmailSendingArgs`. * Also specifies a different connection string for the `EmailSendingArgs`.
`JobQueueConfiguration` class has some additional options in its constructor; `JobQueueConfiguration` class has some additional options in its constructor;
* `queueName`: The queue name that is used for this job. The prefix is not added, so you need to specify the full name of the queue. * `queueName`: The queue name that is used for this job. The prefix is not added, so you need to specify the full name of the queue.
* `DelayedQueueName`: The delayed queue name that is used for delayed execution of job. The prefix is not added, so you need to specify the full name of the queue.
* `connectionName`: The RabbitMQ connection name (see the connection configuration above). This is optional and the default value is `Default`. * `connectionName`: The RabbitMQ connection name (see the connection configuration above). This is optional and the default value is `Default`.
* `durable` (optional, default: `true`). * `durable` (optional, default: `true`).
* `exclusive` (optional, default: `false`). * `exclusive` (optional, default: `false`).
* `autoDelete` (optional, default: `false`) * `autoDelete` (optional, default: `false`).
* `PrefetchCount` (optional, default: null)
See the RabbitMQ documentation if you want to understand the `durable`, `exclusive` and `autoDelete` options better, while most of the times the default configuration is what you want. See the RabbitMQ documentation if you want to understand the `durable`, `exclusive` and `autoDelete` options better, while most of the times the default configuration is what you want.

3
docs/en/Distributed-Event-Bus-RabbitMQ-Integration.md

@ -141,13 +141,14 @@ Configure<AbpRabbitMqOptions>(options =>
}); });
```` ````
**Example: Configure the client and exchange names** **Example: Configure the client, exchange names and prefetchCount**
````csharp ````csharp
Configure<AbpRabbitMqEventBusOptions>(options => Configure<AbpRabbitMqEventBusOptions>(options =>
{ {
options.ClientName = "TestApp1"; options.ClientName = "TestApp1";
options.ExchangeName = "TestMessages"; options.ExchangeName = "TestMessages";
options.PrefetchCount = 1;
}); });
```` ````

10
docs/zh-Hans/Background-Jobs-RabbitMq.md

@ -126,25 +126,31 @@ Configure<AbpRabbitMqOptions>(options =>
Configure<AbpRabbitMqBackgroundJobOptions>(options => Configure<AbpRabbitMqBackgroundJobOptions>(options =>
{ {
options.DefaultQueueNamePrefix = "my_app_jobs."; options.DefaultQueueNamePrefix = "my_app_jobs.";
options.DefaultDelayedQueueNamePrefix = "my_app_jobs.delayed"
options.PrefetchCount = 1;
options.JobQueues[typeof(EmailSendingArgs)] = options.JobQueues[typeof(EmailSendingArgs)] =
new JobQueueConfiguration( new JobQueueConfiguration(
typeof(EmailSendingArgs), typeof(EmailSendingArgs),
queueName: "my_app_jobs.emails", queueName: "my_app_jobs.emails",
connectionName: "SecondConnection" connectionName: "SecondConnection",
delayedQueueName:"my_app_jobs.emails.delayed"
); );
}); });
``` ```
- 这个示例将默认的队列名前缀设置为 `my_app_jobs.`,如果多个项目都使用的同一个 RabbitMQ 服务,设置不同的前缀可以避免执行其他项目的后台作业. - 这个示例将默认的队列名前缀设置为 `my_app_jobs.`并且设置默认的延迟队列名为 `my_app_jobs.delayed`,如果多个项目都使用的同一个 RabbitMQ 服务,设置不同的前缀可以避免执行其他项目的后台作业.
- 设置了预取数量, 用于所有队列.
- 这里还设置了 `EmailSendingArgs` 绑定的 RabbitMQ 连接. - 这里还设置了 `EmailSendingArgs` 绑定的 RabbitMQ 连接.
`JobQueueConfiguration` 类的构造函数中,还有一些其他的可选参数. `JobQueueConfiguration` 类的构造函数中,还有一些其他的可选参数.
- `queueName`: 指定后台作业对应的队列名称(全名). - `queueName`: 指定后台作业对应的队列名称(全名).
* `DelayedQueueName`: 指定后台延迟执行的作业对于的队列名称(全名).
- `connectionName`: 后台作业对应的 RabbitMQ 连接名称,默认是 `Default`. - `connectionName`: 后台作业对应的 RabbitMQ 连接名称,默认是 `Default`.
- `durable`: 可选参数,默认为 `true`. - `durable`: 可选参数,默认为 `true`.
- `exclusive`: 可选参数,默认为 `false`. - `exclusive`: 可选参数,默认为 `false`.
- `autoDelete`: 可选参数,默认为 `false`. - `autoDelete`: 可选参数,默认为 `false`.
* `PrefetchCount` (可选参数, 默认为: null)
如果你想要更多地了解 `durable`,`exclusive`,`autoDelete` 的用法,请阅读 RabbitMQ 提供的文档. 如果你想要更多地了解 `durable`,`exclusive`,`autoDelete` 的用法,请阅读 RabbitMQ 提供的文档.

3
docs/zh-Hans/Distributed-Event-Bus-RabbitMQ-Integration.md

@ -141,13 +141,14 @@ Configure<AbpRabbitMqOptions>(options =>
}); });
```` ````
**示例: 配置客户端和交换机名称** **示例: 配置客户端,交换机名称和预取数量**
````csharp ````csharp
Configure<AbpRabbitMqEventBusOptions>(options => Configure<AbpRabbitMqEventBusOptions>(options =>
{ {
options.ClientName = "TestApp1"; options.ClientName = "TestApp1";
options.ExchangeName = "TestMessages"; options.ExchangeName = "TestMessages";
options.PrefetchCount = 1;
}); });
```` ````

2
framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/AbpRabbitMqBackgroundJobOptions.cs

@ -19,6 +19,8 @@ public class AbpRabbitMqBackgroundJobOptions
/// Default value: "AbpBackgroundJobsDelayed." /// Default value: "AbpBackgroundJobsDelayed."
/// </summary> /// </summary>
public string DefaultDelayedQueueNamePrefix { get; set; } public string DefaultDelayedQueueNamePrefix { get; set; }
public ushort? PrefetchCount { get; set; }
public AbpRabbitMqBackgroundJobOptions() public AbpRabbitMqBackgroundJobOptions()
{ {

10
framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueue.cs

@ -66,7 +66,8 @@ public class JobQueue<TArgs> : IJobQueue<TArgs>
new JobQueueConfiguration( new JobQueueConfiguration(
typeof(TArgs), typeof(TArgs),
AbpRabbitMqBackgroundJobOptions.DefaultQueueNamePrefix + JobConfiguration.JobName, AbpRabbitMqBackgroundJobOptions.DefaultQueueNamePrefix + JobConfiguration.JobName,
AbpRabbitMqBackgroundJobOptions.DefaultDelayedQueueNamePrefix + JobConfiguration.JobName AbpRabbitMqBackgroundJobOptions.DefaultDelayedQueueNamePrefix + JobConfiguration.JobName,
prefetchCount: AbpRabbitMqBackgroundJobOptions.PrefetchCount
); );
} }
@ -140,9 +141,14 @@ public class JobQueue<TArgs> : IJobQueue<TArgs>
if (AbpBackgroundJobOptions.IsJobExecutionEnabled) if (AbpBackgroundJobOptions.IsJobExecutionEnabled)
{ {
if (QueueConfiguration.PrefetchCount.HasValue)
{
ChannelAccessor.Channel.BasicQos(0, QueueConfiguration.PrefetchCount.Value, false);
}
Consumer = new AsyncEventingBasicConsumer(ChannelAccessor.Channel); Consumer = new AsyncEventingBasicConsumer(ChannelAccessor.Channel);
Consumer.Received += MessageReceived; Consumer.Received += MessageReceived;
//TODO: What BasicConsume returns? //TODO: What BasicConsume returns?
ChannelAccessor.Channel.BasicConsume( ChannelAccessor.Channel.BasicConsume(
queue: QueueConfiguration.QueueName, queue: QueueConfiguration.QueueName,

6
framework/src/Volo.Abp.BackgroundJobs.RabbitMQ/Volo/Abp/BackgroundJobs/RabbitMQ/JobQueueConfiguration.cs

@ -20,12 +20,14 @@ public class JobQueueConfiguration : QueueDeclareConfiguration
string connectionName = null, string connectionName = null,
bool durable = true, bool durable = true,
bool exclusive = false, bool exclusive = false,
bool autoDelete = false) bool autoDelete = false,
ushort? prefetchCount = null)
: base( : base(
queueName, queueName,
durable, durable,
exclusive, exclusive,
autoDelete) autoDelete,
prefetchCount)
{ {
JobArgsType = jobArgsType; JobArgsType = jobArgsType;
ConnectionName = connectionName; ConnectionName = connectionName;

2
framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/AbpRabbitMqEventBusOptions.cs

@ -13,6 +13,8 @@ public class AbpRabbitMqEventBusOptions
public string ExchangeName { get; set; } public string ExchangeName { get; set; }
public string ExchangeType { get; set; } public string ExchangeType { get; set; }
public ushort? PrefetchCount { get; set; }
public string GetExchangeTypeOrDefault() public string GetExchangeTypeOrDefault()
{ {

3
framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs

@ -78,7 +78,8 @@ public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDe
AbpRabbitMqEventBusOptions.ClientName, AbpRabbitMqEventBusOptions.ClientName,
durable: true, durable: true,
exclusive: false, exclusive: false,
autoDelete: false autoDelete: false,
prefetchCount: AbpRabbitMqEventBusOptions.PrefetchCount
), ),
AbpRabbitMqEventBusOptions.ConnectionName AbpRabbitMqEventBusOptions.ConnectionName
); );

6
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs

@ -13,6 +13,8 @@ public class QueueDeclareConfiguration
public bool Exclusive { get; set; } public bool Exclusive { get; set; }
public bool AutoDelete { get; set; } public bool AutoDelete { get; set; }
public ushort? PrefetchCount { get; set; }
public IDictionary<string, object> Arguments { get; } public IDictionary<string, object> Arguments { get; }
@ -20,13 +22,15 @@ public class QueueDeclareConfiguration
[NotNull] string queueName, [NotNull] string queueName,
bool durable = true, bool durable = true,
bool exclusive = false, bool exclusive = false,
bool autoDelete = false) bool autoDelete = false,
ushort? prefetchCount = null)
{ {
QueueName = queueName; QueueName = queueName;
Durable = durable; Durable = durable;
Exclusive = exclusive; Exclusive = exclusive;
AutoDelete = autoDelete; AutoDelete = autoDelete;
Arguments = new Dictionary<string, object>(); Arguments = new Dictionary<string, object>();
PrefetchCount = prefetchCount;
} }
public virtual QueueDeclareOk Declare(IModel channel) public virtual QueueDeclareOk Declare(IModel channel)

7
framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs

@ -165,9 +165,14 @@ public class RabbitMqMessageConsumer : IRabbitMqMessageConsumer, ITransientDepen
arguments: Queue.Arguments arguments: Queue.Arguments
); );
if (Queue.PrefetchCount.HasValue)
{
Channel.BasicQos(0, Queue.PrefetchCount.Value, false);
}
var consumer = new AsyncEventingBasicConsumer(Channel); var consumer = new AsyncEventingBasicConsumer(Channel);
consumer.Received += HandleIncomingMessageAsync; consumer.Received += HandleIncomingMessageAsync;
Channel.BasicConsume( Channel.BasicConsume(
queue: Queue.QueueName, queue: Queue.QueueName,
autoAck: false, autoAck: false,

Loading…
Cancel
Save