mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-09-24 14:38:45 +00:00
applied CompleteableFuture pattern to shards. Formatting. Changed
function for getting sharding keys
This commit is contained in:
+81
-80
@@ -30,7 +30,7 @@ import com.sap.sse.landscape.aws.AwsApplicationReplicaSet;
|
||||
import com.sap.sse.landscape.aws.AwsAutoScalingGroup;
|
||||
import com.sap.sse.landscape.aws.AwsLandscape;
|
||||
import com.sap.sse.landscape.aws.AwsShard;
|
||||
import com.sap.sse.landscape.aws.ShardNameDTO;
|
||||
import com.sap.sse.landscape.aws.ShardName;
|
||||
import com.sap.sse.landscape.aws.TargetGroup;
|
||||
|
||||
import software.amazon.awssdk.services.autoscaling.model.AutoScalingGroup;
|
||||
@@ -69,7 +69,7 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
private static final Logger logger = Logger.getLogger(AwsApplicationReplicaSetImpl.class.getName());
|
||||
private static final String ARCHIVE_SERVER_NAME = "ARCHIVE";
|
||||
private static final long serialVersionUID = 6895927683667795173L;
|
||||
private final Map<AwsShard<ShardingKey>, Iterable<ShardingKey>> shards;
|
||||
private final CompletableFuture<Map<AwsShard<ShardingKey>, Iterable<ShardingKey>>> shards;
|
||||
private final CompletableFuture<AwsAutoScalingGroup> autoScalingGroup;
|
||||
private final CompletableFuture<Rule> defaultRedirectRule;
|
||||
private final CompletableFuture<String> hostedZoneId;
|
||||
@@ -84,7 +84,7 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
CompletableFuture<Map<TargetGroup<ShardingKey>, Iterable<TargetHealthDescription>>> allTargetGroupsInRegion,
|
||||
CompletableFuture<Map<Listener, Iterable<Rule>>> allLoadBalancerRulesInRegion,
|
||||
CompletableFuture<Iterable<AutoScalingGroup>> allAutoScalingGroups,
|
||||
CompletableFuture<Iterable<LaunchConfiguration>> allLaunchConfigurations, DNSCache dnsCache, AwsLandscape<ShardingKey> landscape) {
|
||||
CompletableFuture<Iterable<LaunchConfiguration>> allLaunchConfigurations, DNSCache dnsCache) {
|
||||
super(replicaSetAndServerName, hostname, master, replicas);
|
||||
autoScalingGroup = new CompletableFuture<>();
|
||||
defaultRedirectRule = new CompletableFuture<>();
|
||||
@@ -94,13 +94,13 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
masterTargetGroup = new CompletableFuture<>();
|
||||
publicTargetGroup = new CompletableFuture<>();
|
||||
resourceRecordSet = new CompletableFuture<>();
|
||||
shards = new HashMap<>();
|
||||
shards = new CompletableFuture<>();
|
||||
allLoadBalancersInRegion.thenCompose(loadBalancers->
|
||||
allTargetGroupsInRegion.thenCompose(targetGroupsAndTheirTargetHealthDescriptions->
|
||||
allLoadBalancerRulesInRegion.thenCompose(listenersAndTheirRules->
|
||||
allAutoScalingGroups.thenCompose(autoScalingGroups->
|
||||
allLaunchConfigurations.handle((launchConfigurations, e)->establishState(
|
||||
loadBalancers, targetGroupsAndTheirTargetHealthDescriptions, listenersAndTheirRules, autoScalingGroups, launchConfigurations, dnsCache, landscape))))))
|
||||
loadBalancers, targetGroupsAndTheirTargetHealthDescriptions, listenersAndTheirRules, autoScalingGroups, launchConfigurations, dnsCache))))))
|
||||
.handle((v, e)->{
|
||||
if (e != null) {
|
||||
logger.log(Level.SEVERE, "Exception while trying to establish state of application replica set "+getName(), e);
|
||||
@@ -109,8 +109,8 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
});
|
||||
}
|
||||
|
||||
public Map<AwsShard<ShardingKey>, Iterable<ShardingKey>> getShards(){
|
||||
return shards;
|
||||
public Map<AwsShard<ShardingKey>, Iterable<ShardingKey>> getShards() throws Exception{
|
||||
return shards.get();
|
||||
}
|
||||
|
||||
public AwsApplicationReplicaSetImpl(String replicaSetAndServerName, ProcessT master,
|
||||
@@ -122,7 +122,7 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
CompletableFuture<Iterable<LaunchConfiguration>> allLaunchConfigurations, DNSCache dnsCache) {
|
||||
this(replicaSetAndServerName, /* hostname to be inferred */ null, master, replicas, allLoadBalancersInRegion,
|
||||
allTargetGroupsInRegion, allLoadBalancerRulesInRegion, allAutoScalingGroups, allLaunchConfigurations,
|
||||
dnsCache, landscape);
|
||||
dnsCache);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -163,7 +163,7 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
private Void establishState(Iterable<ApplicationLoadBalancer<ShardingKey>> loadBalancers,
|
||||
Map<TargetGroup<ShardingKey>, Iterable<TargetHealthDescription>> targetGroupsAndTheirTargetHealthDescriptions,
|
||||
Map<Listener, Iterable<Rule>> listenersAndTheirRules, Iterable<AutoScalingGroup> autoScalingGroups,
|
||||
Iterable<LaunchConfiguration> launchConfigurations, DNSCache dnsCache, AwsLandscape<ShardingKey> landscape) {
|
||||
Iterable<LaunchConfiguration> launchConfigurations, DNSCache dnsCache) {
|
||||
TargetGroup<ShardingKey> myMasterTargetGroup = null;
|
||||
TargetGroup<ShardingKey> singleTargetGroupCandidate = null;
|
||||
for (final Entry<TargetGroup<ShardingKey>, Iterable<TargetHealthDescription>> e : targetGroupsAndTheirTargetHealthDescriptions.entrySet()) {
|
||||
@@ -198,7 +198,9 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
}
|
||||
}
|
||||
}
|
||||
shards.putAll(establishShards(targetGroupsAndTheirTargetHealthDescriptions, listenersAndTheirRules,autoScalingGroups,launchConfigurations, landscape));
|
||||
if (!shards.isDone()) {
|
||||
shards.complete(establishShards(targetGroupsAndTheirTargetHealthDescriptions, listenersAndTheirRules,autoScalingGroups,launchConfigurations));
|
||||
}
|
||||
if (!autoScalingGroup.isDone()) { // no auto-scaling group was found after having looked at all target groups
|
||||
autoScalingGroup.complete(null);
|
||||
}
|
||||
@@ -310,10 +312,10 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
HashMap<AwsShard<ShardingKey>, Iterable<ShardingKey>> establishShards(
|
||||
private HashMap<AwsShard<ShardingKey>, Iterable<ShardingKey>> establishShards(
|
||||
Map<TargetGroup<ShardingKey>, Iterable<TargetHealthDescription>> targetGroupsAndTheirTargetHealthDescriptions,
|
||||
Map<Listener, Iterable<Rule>> listenersAndTheirRules, Iterable<AutoScalingGroup> autoscalinggroups,
|
||||
Iterable<LaunchConfiguration> launchConfigurations, AwsLandscape<ShardingKey> landscape) {
|
||||
Map<Listener, Iterable<Rule>> listenersAndTheirRules, Iterable<AutoScalingGroup> autoScalingGroups,
|
||||
Iterable<LaunchConfiguration> launchConfigurations) {
|
||||
HashMap<AwsShard<ShardingKey>, Iterable<ShardingKey>> shardMap = new HashMap<>();
|
||||
for (final Entry<TargetGroup<ShardingKey>, Iterable<TargetHealthDescription>> e : targetGroupsAndTheirTargetHealthDescriptions
|
||||
.entrySet()) {
|
||||
@@ -322,82 +324,46 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
&& e.getKey().getHealthCheckPort() == getMaster().getPort()) {
|
||||
// an HTTP(S) target group that is currently active in a load balancer and forwards to this replica
|
||||
// set master's port
|
||||
Set<Rule> pathRules = getListenerRuleWithPathAnToReplica(listenersAndTheirRules, e.getKey());
|
||||
final Set<Rule> pathRules = getListenerRulesWithPathAnToReplica(listenersAndTheirRules, e.getKey());
|
||||
if (!pathRules.isEmpty() && isShardName(e.getKey().getName())) {
|
||||
// Is shard
|
||||
ArrayList<ShardingKey> keys = new ArrayList<ShardingKey>();
|
||||
// get every shardingkey from the rules
|
||||
for (Rule r : pathRules) {
|
||||
for (RuleCondition condition : r.conditions()) {
|
||||
|
||||
if (condition.pathPatternConfig() != null) {
|
||||
for (String s : condition.pathPatternConfig().values()) {
|
||||
keys.add((ShardingKey) s.toString());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
String tagname = null;
|
||||
Iterable<TagDescription> tagsDesc = landscape.getTargetGroupTags(e.getKey().getTargetGroupArn(),
|
||||
e.getKey().getRegion());
|
||||
Set<String> keys = getShardingKeys(listenersAndTheirRules, e.getKey());
|
||||
String tagName = null;
|
||||
Iterable<TagDescription> tagsDesc = e.getKey().getTagDescriptions();
|
||||
for (TagDescription des : tagsDesc) {
|
||||
Iterable<Tag> tag = Util.filter(des.tags(), t -> t.key().equals(ShardNameDTO.TAG_KEY));
|
||||
Iterable<Tag> tag = Util.filter(des.tags(), t -> t.key().equals(ShardName.TAG_KEY));
|
||||
if (!Util.isEmpty(tag)) {
|
||||
tagname = tag.iterator().next().value();
|
||||
tagName = tag.iterator().next().value();
|
||||
break;
|
||||
}
|
||||
}
|
||||
ShardNameDTO shardname;
|
||||
ShardName shardName;
|
||||
try {
|
||||
shardname = ShardNameDTO.parse(e.getKey().getName(), tagname);
|
||||
shardName = ShardName.parse(e.getKey().getName(), tagName);
|
||||
} catch (Exception e1) {
|
||||
logger.info(e1.getMessage());
|
||||
// This entry is no valid shard
|
||||
continue;
|
||||
}
|
||||
AwsShardImpl<ShardingKey> shard = new AwsShardImpl<ShardingKey>(getName(), keys, e.getKey(),
|
||||
shardname, e.getKey().getLoadBalancer(), pathRules);
|
||||
AwsShardImpl<ShardingKey> shard = new AwsShardImpl<ShardingKey>(getName(),
|
||||
Util.asList(Util.map(keys, s -> (ShardingKey) s)), e.getKey(), shardName,
|
||||
e.getKey().getLoadBalancer(), pathRules);
|
||||
shardMap.put(shard, shard.getKeys());
|
||||
shard.setAutoscalingGroup(
|
||||
getShardAutoscalinggroup(e.getKey(), autoscalinggroups, launchConfigurations));
|
||||
|
||||
getShardAutoscalinggroup(e.getKey(), autoScalingGroups, launchConfigurations));
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
return shardMap;
|
||||
};
|
||||
|
||||
private boolean isShardName(String requestedName) {
|
||||
|
||||
return ShardNameDTO.isValidTargetGroupName(requestedName);
|
||||
return ShardName.isValidTargetGroupName(requestedName);
|
||||
}
|
||||
|
||||
public ShardNameDTO getNextShardName(String name) throws Exception{
|
||||
return ShardNameDTO.create(getName(), name);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param name
|
||||
* string where the shardname should be extracted
|
||||
* Naming-pattern: {prefix-}{getName()}-{$shardname$}-$shardCounter$>
|
||||
* @return
|
||||
* shardname
|
||||
*/
|
||||
public String getShardName(String name) {
|
||||
final int firstHyphen, secondHyphen;
|
||||
final String ret;
|
||||
secondHyphen =name.lastIndexOf('-');
|
||||
firstHyphen =name.lastIndexOf('-', secondHyphen-1);
|
||||
if(firstHyphen < 3) {
|
||||
//If there was no second '-' or if '-' was in the first 3 indize because this could be a prefix -> there was no shardname given
|
||||
ret = null;
|
||||
}else {
|
||||
ret = name.substring(firstHyphen+1, secondHyphen);
|
||||
}
|
||||
return ret;
|
||||
|
||||
|
||||
@Override
|
||||
public ShardName getNewShardName(String shardName) throws Exception{
|
||||
return ShardName.create(getName(), shardName);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -437,16 +403,16 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
}
|
||||
|
||||
private boolean hasPathCondition(Map<Listener, Iterable<Rule>> listenersAndTheirRules,
|
||||
TargetGroup<ShardingKey> publicTargetGroupCandidate) {
|
||||
final String publicTargetGroupCandidateArn = publicTargetGroupCandidate.getTargetGroupArn();
|
||||
TargetGroup<ShardingKey> shardingTargetGroupCandidate) {
|
||||
final String shardTargetGroupCandidateArn = shardingTargetGroupCandidate.getTargetGroupArn();
|
||||
for (final Entry<Listener, Iterable<Rule>> e : listenersAndTheirRules.entrySet()) {
|
||||
if (Util.equalsWithNull(e.getKey().loadBalancerArn(), publicTargetGroupCandidate.getLoadBalancerArn())) {
|
||||
if (Util.equalsWithNull(e.getKey().loadBalancerArn(), shardingTargetGroupCandidate.getLoadBalancerArn())) {
|
||||
for (final Rule rule : e.getValue()) {
|
||||
for (final Action action : rule.actions()) {
|
||||
if (action.type() == ActionTypeEnum.FORWARD) {
|
||||
for (final TargetGroupTuple targetGroupTuple : action.forwardConfig().targetGroups()) {
|
||||
if (Util.equalsWithNull(targetGroupTuple.targetGroupArn(),
|
||||
publicTargetGroupCandidateArn)) {
|
||||
shardTargetGroupCandidateArn)) {
|
||||
for (final RuleCondition condition : rule.conditions()) {
|
||||
if (Util.equalsWithNull(condition.field(), "path-pattern")) {
|
||||
return true;
|
||||
@@ -488,9 +454,10 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
return false;
|
||||
}
|
||||
|
||||
private Set<Rule> getListenerRuleWithPathAnToReplica(Map<Listener, Iterable<Rule>> listenersAndTheirRules, TargetGroup<ShardingKey> shardTargetGroupCandidate) {
|
||||
private Set<Rule> getListenerRulesWithPathAnToReplica(Map<Listener, Iterable<Rule>> listenersAndTheirRules,
|
||||
TargetGroup<ShardingKey> shardTargetGroupCandidate) {
|
||||
final String publicTargetGroupCandidateArn = shardTargetGroupCandidate.getTargetGroupArn();
|
||||
Set<Rule> res =new HashSet<Rule>();
|
||||
Set<Rule> res = new HashSet<Rule>();
|
||||
for (final Entry<Listener, Iterable<Rule>> e : listenersAndTheirRules.entrySet()) {
|
||||
if (Util.equalsWithNull(e.getKey().loadBalancerArn(), shardTargetGroupCandidate.getLoadBalancerArn())) {
|
||||
for (final Rule rule : e.getValue()) {
|
||||
@@ -501,19 +468,53 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
publicTargetGroupCandidateArn)) {
|
||||
for (final RuleCondition condition : rule.conditions()) {
|
||||
if (Util.equalsWithNull(condition.field(), "path-pattern")) {
|
||||
for (final RuleCondition condition2 : rule.conditions()) {
|
||||
if (Util.equalsWithNull(condition2.field(), "http-header")
|
||||
for (final RuleCondition ruleCondition : rule.conditions()) {
|
||||
if (Util.equalsWithNull(ruleCondition.field(), "http-header")
|
||||
&& Util.equalsWithNull(
|
||||
condition2.httpHeaderConfig().httpHeaderName(),
|
||||
HttpRequestHeaderConstants.HEADER_KEY_FORWARD_TO)
|
||||
&& condition2.httpHeaderConfig().values().contains(
|
||||
ruleCondition.httpHeaderConfig().httpHeaderName(),
|
||||
HttpRequestHeaderConstants.HEADER_FORWARD_TO_REPLICA
|
||||
.getA())
|
||||
&& ruleCondition.httpHeaderConfig().values().contains(
|
||||
HttpRequestHeaderConstants.HEADER_FORWARD_TO_REPLICA
|
||||
.getB())) {
|
||||
res.add(rule);
|
||||
//return true is conditions contains a path and header forward to replica
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return res;
|
||||
}
|
||||
|
||||
private Set<String> getShardingKeys(Map<Listener, Iterable<Rule>> listenersAndTheirRules,
|
||||
TargetGroup<ShardingKey> shardTargetGroupCandidate) {
|
||||
final String publicTargetGroupCandidateArn = shardTargetGroupCandidate.getTargetGroupArn();
|
||||
Set<String> res = new HashSet<>();
|
||||
for (final Entry<Listener, Iterable<Rule>> e : listenersAndTheirRules.entrySet()) {
|
||||
if (Util.equalsWithNull(e.getKey().loadBalancerArn(), shardTargetGroupCandidate.getLoadBalancerArn())) {
|
||||
for (final Rule rule : e.getValue()) {
|
||||
for (final Action action : rule.actions()) {
|
||||
if (action.type() == ActionTypeEnum.FORWARD) {
|
||||
for (final TargetGroupTuple targetGroupTuple : action.forwardConfig().targetGroups()) {
|
||||
if (Util.equalsWithNull(targetGroupTuple.targetGroupArn(),
|
||||
publicTargetGroupCandidateArn)) {
|
||||
for (final RuleCondition condition : rule.conditions()) {
|
||||
if (Util.equalsWithNull(condition.field(), "http-header")
|
||||
&& Util.equalsWithNull(condition.httpHeaderConfig().httpHeaderName(),
|
||||
HttpRequestHeaderConstants.HEADER_FORWARD_TO_REPLICA.getA())
|
||||
&& condition.httpHeaderConfig().values().contains(
|
||||
HttpRequestHeaderConstants.HEADER_FORWARD_TO_REPLICA.getB())) {
|
||||
for (final RuleCondition ruleCondition : rule.conditions()) {
|
||||
if (Util.equalsWithNull(ruleCondition.field(), "path-pattern")) {
|
||||
res.addAll(ruleCondition.values());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -607,10 +608,10 @@ implements AwsApplicationReplicaSet<ShardingKey, MetricsT, ProcessT> {
|
||||
targetGroups.add(shard.getTargetGroup());
|
||||
// remove rules for targetgrouo
|
||||
landscape.deleteLoadBalancerListenerRules(shard.getTargetGroup().getRegion(),
|
||||
getLoadBalancer().getRulesForTargetGroups(targetGroups).toArray(new Rule[1]));
|
||||
Util.toArray(getLoadBalancer().getRulesForTargetGroups(targetGroups), new Rule[0]));
|
||||
// remove autoscaling group
|
||||
landscape.removeAutoScalingGroup(shard.getAutoScalingGroup());
|
||||
// remove targetgroup
|
||||
landscape.deleteTargetGroup(shard.getTargetGroup());
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user