Commit 0c831314 authored by Savorboard's avatar Savorboard

add connection pool size config to KafkaOption.

parent 2e81f2bc
using System; using System;
using System.Collections.Concurrent;
using System.Collections.Generic; using System.Collections.Generic;
using System.Linq; using System.Linq;
...@@ -16,16 +17,21 @@ namespace DotNetCore.CAP ...@@ -16,16 +17,21 @@ namespace DotNetCore.CAP
/// Topic configuration parameters are specified via the "default.topic.config" sub-dictionary config parameter. /// Topic configuration parameters are specified via the "default.topic.config" sub-dictionary config parameter.
/// </para> /// </para>
/// </summary> /// </summary>
public readonly IDictionary<string, object> MainConfig; public readonly ConcurrentDictionary<string, object> MainConfig;
private IEnumerable<KeyValuePair<string, object>> _kafkaConfig; private IEnumerable<KeyValuePair<string, object>> _kafkaConfig;
public KafkaOptions() public KafkaOptions()
{ {
MainConfig = new Dictionary<string, object>(); MainConfig = new ConcurrentDictionary<string, object>();
} }
/// <summary>
/// Producer connection pool size, default is 10
/// </summary>
public int ConnectionPoolSize { get; set; } = 10;
/// <summary> /// <summary>
/// The `bootstrap.servers` item config of <see cref="MainConfig" />. /// The `bootstrap.servers` item config of <see cref="MainConfig" />.
/// <para> /// <para>
......
...@@ -8,10 +8,7 @@ namespace DotNetCore.CAP.Kafka ...@@ -8,10 +8,7 @@ namespace DotNetCore.CAP.Kafka
{ {
public class ConnectionPool : IConnectionPool, IDisposable public class ConnectionPool : IConnectionPool, IDisposable
{ {
private const int DefaultPoolSize = 15;
private readonly Func<Producer> _activator; private readonly Func<Producer> _activator;
private readonly ConcurrentQueue<Producer> _pool = new ConcurrentQueue<Producer>(); private readonly ConcurrentQueue<Producer> _pool = new ConcurrentQueue<Producer>();
private int _count; private int _count;
...@@ -19,8 +16,7 @@ namespace DotNetCore.CAP.Kafka ...@@ -19,8 +16,7 @@ namespace DotNetCore.CAP.Kafka
public ConnectionPool(KafkaOptions options) public ConnectionPool(KafkaOptions options)
{ {
_maxSize = DefaultPoolSize; _maxSize = options.ConnectionPoolSize;
_activator = CreateActivator(options); _activator = CreateActivator(options);
} }
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment