【问题标题】:Cluster sharding client not connecting with host集群分片客户端未与主机连接
【发布时间】:2017-10-05 09:38:47
【问题描述】:

在最近的调查和Stack over flow question 之后,我意识到集群分片比集群一致的哈希路由器更好。但我无法让 2 进程集群运行。

一个进程是种子,另一个是客户端。 Seed 节点似乎不断地抛出死信消息(见本问题的结尾)。

此种子 HOCON 如下:

akka {
loglevel = "INFO"                    

actor {
    provider = "Akka.Cluster.ClusterActorRefProvider, Akka.Cluster"
    serializers {
        wire = "Akka.Serialization.WireSerializer, Akka.Serialization.Wire"
    }
    serialization-bindings {
        "System.Object" = wire
    }
}                    

remote {
    dot-netty.tcp {
        hostname = "127.0.0.1"
        port = 5000
    }
}

persistence {
    journal {
        plugin = "akka.persistence.journal.sql-server"
        sql-server {
            class = "Akka.Persistence.SqlServer.Journal.SqlServerJournal, Akka.Persistence.SqlServer"
            schema-name = dbo
            auto-initialize = on
            connection-string = "Data Source=localhost;Integrated Security=True;MultipleActiveResultSets=True;Initial Catalog=ClusterExperiment01"
            plugin-dispatcher = "akka.actor.default- dispatcher"
            connection-timeout = 30s
            table-name = EventJournal
            timestamp-provider = "Akka.Persistence.Sql.Common.Journal.DefaultTimestampProvider, Akka.Persistence.Sql.Common"
            metadata-table-name = Metadata
        }
    }

    sharding {
        connection-string = "Data Source=localhost;Integrated Security=True;MultipleActiveResultSets=True;Initial Catalog=ClusterExperiment01"
        auto-initialize = on
        plugin-dispatcher = "akka.actor.default-dispatcher"
        class = "Akka.Persistence.SqlServer.Journal.SqlServerJournal, Akka.Persistence.SqlServer"
        connection-timeout = 30s
        schema-name = dbo
        table-name = ShardingJournal
        timestamp-provider = "Akka.Persistence.Sql.Common.Journal.DefaultTimestampProvider, Akka.Persistence.Sql.Common"
        metadata-table-name = ShardingMetadata
    }
}

snapshot-store {
    sharding {
        class = "Akka.Persistence.SqlServer.Snapshot.SqlServerSnapshotStore, Akka.Persistence.SqlServer"
        plugin-dispatcher = "akka.actor.default-dispatcher"
        connection-string = "Data Source=localhost;Integrated Security=True;MultipleActiveResultSets=True;Initial Catalog=ClusterExperiment01"
        connection-timeout = 30s
        schema-name = dbo
        table-name = ShardingSnapshotStore
        auto-initialize = on
    }
}

cluster {
    seed-nodes = ["akka.tcp://my-cluster-system@127.0.0.1:5000"]
    roles = ["Seed"]

    sharding {
        journal-plugin-id = "akka.persistence.sharding"
        snapshot-plugin-id = "akka.snapshot-store.sharding"
    }
}}

我有一个方法可以将上面的内容变成这样的配置:

var config = NodeConfig.Create(/* HOCON above */).WithFallback(ClusterSingletonManager.DefaultConfig());

如果没有“WithFallback”,我会从配置生成中得到一个空引用异常。

然后像这样生成系统:

var system = ActorSystem.Create("my-cluster-system", config);

客户端以相同的方式创建其系统,HOCON 几乎相同,除了:

{
remote {
    dot-netty.tcp {
        hostname = "127.0.0.1"
        port = 5001
    }
}
cluster {
    seed-nodes = ["akka.tcp://my-cluster-system@127.0.0.1:5000"]
    roles = ["Client"]
    role.["Seed"].min-nr-of-members = 1
    sharding {
        journal-plugin-id = "akka.persistence.sharding"
        snapshot-plugin-id = "akka.snapshot-store.sharding"
    }
}}

种子节点像这样创建分片:

ClusterSharding.Get(system).Start(
   typeName: "company-router",
   entityProps: Props.Create(() => new CompanyDeliveryActor()),                    
   settings: ClusterShardingSettings.Create(system),
   messageExtractor: new RouteExtractor(100)
);

客户端创建一个分片代理,如下所示:

ClusterSharding.Get(system).StartProxy(
    typeName: "company-router",
    role: "Seed",
    messageExtractor: new RouteExtractor(100));

RouteExtractor 是:

public class RouteExtractor : HashCodeMessageExtractor
{
    public RouteExtractor(int maxNumberOfShards) : base(maxNumberOfShards)
    {   
    }
    public override string EntityId(object message) => (message as IHasRouting)?.Company?.VolumeId.ToString();
    public override object EntityMessage(object message) => message;
}

在这种情况下,VolumeId 始终相同(仅用于实验)。

两个进程都恢复了,但种子不断向日志抛出此错误:

[INFO][7/05/2017 上午 9:00:58][线程 0003][akka://my-cluster-system/user/sharding /company-routerCoordinator/singleton/coordinator] 来自 akka.tcp 的消息注册 ://my-cluster-system@127.0.0.1:5000/user/sharding/company-router 到 akka://my-cl uster-system/user/sharding/company-routerCoordinator/singleton/coordinator 是 n 未交付。遇到 4 个死信。

附言。我没有使用 Lighthouse。

【问题讨论】:

    标签: akka.net consistent-hashing akka.net-cluster


    【解决方案1】:

    感谢Horusiath,已解决:

    return sharding.Start(
       typeName: "company-router",
       entityProps: Props.Create(() => new CompanyDeliveryActor()),                    
       settings: ClusterShardingSettings.Create(system).WithRole("Seed"),
                    messageExtractor: new RouteExtractor(100)                
                );
    

    集群分片现在正在两个进程之间进行通信。非常感谢您。

    【讨论】:

      【解决方案2】:

      快速浏览一下,您将在客户端节点上启动集群分片代理,并告诉它分片节点是那些使用 seed 角色的节点。当您没有指定任何角色时,这与种子节点上的集群分片定义不匹配。

      由于没有角色可以限制它,种子节点上的集群分片将把集群中的所有节点视为完全能够托管分片参与者 - 包括客户端节点,它没有实例化集群分片(非代理)就可以了。

      这可能不是唯一的问题,但您可以在所有节点上托管集群分片,或者使用 ClusterShardingSettings.Create(system).WithRole("seed") 将您的分片限制在特定的节点子集(具有 seed角色)在集群中。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2017-11-15
        • 1970-01-01
        • 1970-01-01
        • 2014-11-29
        • 2021-09-08
        • 1970-01-01
        • 1970-01-01
        • 2014-08-13
        相关资源
        最近更新 更多