kafka消費者客戶端

Kafka消費者

1.1 消費者與消費者組

消費者與消費者組之間的關係

​ 每一個消費者都隸屬於某一個消費者組,一個消費者組可以包含一個或多個消費者,每一條消息只會被消費者組中的某一個消費者所消費。不同消費者組之間消息的消費是互不干擾的。

為什麼會有消費者組的概念

​ 消費者組出現主要是出於兩個目的:

​ (1) 使整體的消費能力具備橫向的伸縮性。可以適當增加消費者組中消費者的數量,來提高整體的消費能力。但是每一個分區至多被消費者組的中一個消費者所消費,因此當消費者組中消費者數量超過分區數時,多出的消費者不會分配到任何一個分區。當然這是默認的分區分配策略,可通過partition.assignment.strategy進行配置。

​ (2) 實現消息消費的隔離。不同消費者組之間消息消費互不干擾,從而實現發布訂閱這種消息投遞模式。

注意:

​ 消費者隸屬的消費者組可以通過group.id進行配置。消費者組是一個邏輯上的概念,但消費者並不是一個邏輯上的概念,它可以是一個線程,也可以是一個進程。同一個消費者組內的消費者可以部署在同一台機器上,也可以部署在不同的機器上。

1.2 消費者客戶端開發

​ 一個正常的消費邏輯需要具備以下幾個步驟:

  • 配置消費者客戶端參數及創建相應的消費者實例。

  • 訂閱主題

  • 拉取消息並消費

  • 提交消費位移

  • 關閉消費者實例

  public class KafkaConsumerAnalysis {
  public static final String brokerList="node112:9092,node113:9092,node114:9092";
  public static final String topic = "topic-demo";
  public static final String groupId = "group.demo";
  public static final AtomicBoolean isRunning = new AtomicBoolean(true);

  public static Properties initConfig() {
      Properties prop = new Properties();
      prop.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
      prop.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
      prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
      prop.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
      prop.put(ConsumerConfig.CLIENT_DNS_LOOKUP_CONFIG, "consumer.client.di.demo");
      return prop;
  }


  public static void main(String[] args) {
      KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(initConfig());
       
      for (ConsumerRecord<String, String> record : records) {
          System.out.println("topic = " + record.topic() + ", partition =" +                                               record.partition() + ", offset = " + record.offset());
          System.out.println("key = " + record.key() + ", value = " + record.value());
              }
          }
      } catch (Exception e) {
          e.printStackTrace();
      }finally {
          consumer.close();
      }
  }
}

1.2.1 訂閱主題和分區

​ 先來說一下消費者訂閱消息的粒度:一個消費者可以訂閱一個主題、多個主題、或者多個主題的特定分區。主要通過subsribe和assign兩個方法實現訂閱。

(1)訂閱一個主題:

​ public void subscribe(Collection<String> topics),當集合中有一個主題時。

(2)訂閱多個主題:

​ public void subscribe(Collection<String> topics),當集合中有多個主題時。

​ public void subscribe(Pattern pattern),通過正則表達式實現消費者主題的匹配。通過這種方式,如果在消息消費的過程中,又添加了新的能夠匹配到正則的主題,那麼消費者就可以消費到新添加的主題。 consumer.subscribe(Pattern.compile(“topic-.*”));

(3)多個主題的特定分區

​ public void assign(Collection<TopicPartition> partitions),可以實現訂閱某些特定的主題分區。TopicPartition包括兩個屬性:topic(String)和partition(int)。

​ 如果事先不知道有多少分區該如何處理,KafkaConsumer中的partitionFor方法可以獲得指定主題分區的元數據信息:

​ public List<PartitionInfo> partitionsFor(String topic)

​ PartitionInfo的屬性如下:

  
public class PartitionInfo {
  private final String topic;//主題
  private final int partition;//分區
  private final Node leader;//分區leader
  private final Node[] replicas;//分區的AR
  private final Node[] inSyncReplicas;//分區的ISR
  private final Node[] offlineReplicas;//分區的OSR
}

​ 因此也可以通過這個方法實現某個主題的全部訂閱。

​ 需要指出的是,subscribe(Collection)、subscirbe(Pattern)、assign(Collection)方法分別代表了三種不同的訂閱狀態:AUTO_TOPICS、AUTO_PATTREN和USER_ASSIGN,這三種方式是互斥的,消費者只能使用其中一種,否則會報出IllegalStateException。

​ subscirbe方法可以實現消費者自動再平衡的功能。多個消費者的情況下,可以根據分區分配策略自動分配消費者和分區的關係,當消費者增加或減少時,也能實現負載均衡和故障轉移。

​ 如何實現取消訂閱:

​ consumer.unsubscribe()

1.2.2 反序列化

​ KafkaProducer端生產消息進行序列化,同樣消費者就要進行相應的反序列化。相當於根據定義的序列化格式的一個逆序提取數據的過程。

  
import com.gdy.kafka.producer.Company;
import org.apache.kafka.common.errors.SerializationException;
import org.apache.kafka.common.serialization.Deserializer;

import java.io.UnsupportedEncodingException;
import java.nio.ByteBuffer;
import java.util.Map;

public class CompanyDeserializer implements Deserializer<Company> {
  @Override
  public void configure(Map<String, ?> configs, boolean isKey) {

  }

  @Override
  public Company deserialize(String topic, byte[] data) {
      if(data == null) {
          return null;
      }

      if(data.length < 8) {
          throw new SerializationException("size of data received by Deserializer is shorter than expected");
      }

      ByteBuffer buffer = ByteBuffer.wrap(data);
      int nameLength = buffer.getInt();
      byte[] nameBytes = new byte[nameLength];
      buffer.get(nameBytes);
      int addressLen = buffer.getInt();
      byte[] addressBytes = new byte[addressLen];
      buffer.get(addressBytes);
      String name,address;
      try {
          name = new String(nameBytes,"UTF-8");
          address = new String(addressBytes,"UTF-8");
      }catch (UnsupportedEncodingException e) {
          throw new SerializationException("Error accur when deserializing");
      }

      return new Company(name, address);
  }

  @Override
  public void close() {

  }
}

​ 實際生產中需要自定義序列化器和反序列化器時,推薦使用Avro、JSON、Thrift、ProtoBuf或者Protostuff等通用的序列化工具來包裝。

1.2.3 消息消費

​ Kafka中消息的消費是基於拉模式的,kafka消息的消費是一個不斷輪旋的過程,消費者需要做的就是重複的調用poll方法。

  
public ConsumerRecords<K, V> poll(final Duration timeout)

​ 這個方法需要注意的是,如果消費者的緩衝區中有可用的數據,則會立即返回,否則會阻塞至timeout。如果在阻塞時間內緩衝區仍沒有數據,則返回一個空的消息集。timeout的設置取決於應用程序對效應速度的要求。如果應用線程的位移工作是從Kafka中拉取數據並進行消費可以將這個參數設置為Long.MAX_VALUE。

​ 每次poll都會返回一個ConsumerRecords對象,它是ConsumerRecord的集合。對於ConsumerRecord相比於ProducerRecord多了一些屬性:

  
private final String topic;//主題
  private final int partition;//分區
  private final long offset;//偏移量
  private final long timestamp;//時間戳
  private final TimestampType timestampType;//時間戳類型
  private final int serializedKeySize;//序列化key的大小
  private final int serializedValueSize;//序列化value的大小
  private final Headers headers;//headers
  private final K key;//key
  private final V value;//value
  private volatile Long checksum;//CRC32校驗和

​ 另外我們可以按照分區維度對消息進行消費,通過ConsumerRecords.records(TopicPartiton)方法實現。

  
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
    Set<TopicPartition> partitions = records.partitions();
    for (TopicPartition tp : partitions) {
        for (ConsumerRecord<String, String> record : records.records(tp)) {
              System.out.println(record.partition() + " ," + record.value());
        }
    }

​ 另外還可以按照主題維度對消息進行消費,通過ConsumerRecords.records(Topic)實現。

  
for (String topic : topicList) {
      for (ConsumerRecord<String, String> record : records.records(topic)) {
              System.out.println(record.partition() + " ," + record.value());
      }
}

1.2.4 消費者位移提交

​ 首先要 明白一點,消費者位移是要做持久化處理的,否則當發生消費者崩潰或者消費者重平衡時,消費者消費位移無法獲得。舊消費者客戶端是將位移提交到zookeeper上,新消費者客戶端將位移存儲在Kafka內部主題_consumer_offsets中。

​ KafkaConsumer提供了兩個方法position(TopicPatition)和commited(TopicPartition)。

​ public long position(TopicPartition partition)—–獲得下一次拉取數據的偏移量

​ public OffsetAndMetadata committed(TopicPartition partition)—–給定分區的最後一次提交的偏移量。

還有一個概念稱之為lastConsumedOffset,這個指的是最後一次消費的偏移量。

​ 在kafka提交方式有兩種:自動提交和手動提交。

(1)自動位移提交

​ kafka默認情況下採用自動提交,enable.auto.commit的默認值為true。當然自動提交並不是沒消費一次消息就進行提交,而是定期提交,這個定期的周期時間由auto.commit.intervals.ms參數進行配置,默認值為5s,當然這個參數生效的前提就是開啟自動提交。

​ 自動提交會造成重複消費和消息丟失的情況。重複消費很容易理解,因為自動提交實際是延遲提交,因此很容易造成重複消費,然後消息丟失是怎麼產生的?

(2)手動位移提交

​ 開始手動提交的需要配置enable.auto.commit=false。手動提交消費者偏移量,又可分為同步提交和異步提交。

​ 同步提交:

​ 同步提交很簡單,調用commitSync() 方法:

  
while (isRunning.get()) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
        for (ConsumerRecord<String, String> record : records) {
            //consume message
            consumer.commitSync();
        }
}

​ 這樣,每消費一條消息,提交一個偏移量。當然可用過緩存消息的方式,實現批量處理+批量提交:

  
while (isRunning.get()) {
      ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
      for (ConsumerRecord<String, String> record : records) {
            buffer.add(record);
      }
      if (buffer.size() >= minBaches) {
          for (ConsumerRecord<String, String> record : records) {
              //consume message
          }
          consumer.commitSync();
          buffer.clear();
      }
}

​ 還可以通過public void commitSync(final Map<TopicPartition, OffsetAndMetadata> offsets)這個方法實現按照分區粒度進行同步提交。

  
while (isRunning.get()) {
  ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
  for (TopicPartition tp : records.partitions()) {
      List<ConsumerRecord<String, String>> partitionRecords = records.records(tp);
      for (ConsumerRecord record : partitionRecords) {
          //consume message
      }
      long lastConsumerOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
      consumer.commitSync(Collections.singletonMap(tp,new                                                               OffsetAndMetadata(lastConsumerOffset+1)));
  }
}

​ 異步提交:

​ commitAsync異步提交的時候消費者線程不會被阻塞,即可能在提交偏移量的結果還未返回之前,就開始了新一次的拉取數據操作。異步提交可以提升消費者的性能。commitAsync有三個重載:

​ public void commitAsync()

​ public void commitAsync(OffsetCommitCallback callback)

​ public void commitAsync(final Map<TopicPartition, OffsetAndMetadata> offsets, OffsetCommitCallback )

​ 對照同步提交的方法參數,多了一個Callback回調參數,它提供了一個異步提交的回調方法,當消費者位移提交完成后回調OffsetCommitCallback的onComplement方法。以第二個方法為例:

  
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
  for (ConsumerRecord<String, String> record : records) {
      //consume message
  }
  consumer.commitAsync(new OffsetCommitCallback() {
      @Override
      public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception e) {
          if (e == null) {
              System.out.println(offsets);
          }else {
                e.printStackTrace();
          }
      }
});

1.2.5 控制和關閉消費

​ kafkaConsumer提供了pause()和resume() 方法分別實現暫停某些分區在拉取操作時返回數據給客戶端和恢復某些分區向客戶端返回數據的操作:

​ public void pause(Collection<TopicPartition> partitions)

​ public void resume(Collection<TopicPartition> partitions)

​ 優雅停止KafkaConsumer退出消費者循環的方式:

​ (1)不要使用while(true),而是使用while(isRunning.get()),isRunning是一個AtomicBoolean類型,可以在其他地方調用isRunning.set(false)方法退出循環。

​ (2)調用consumer.wakup()方法,wakeup方法是KafkaConsumer中唯一一個可以從其他線程里安全調用的方法,會拋出WakeupException,我們不需要處理這個異常。

​ 跳出循環后一定要显示的執行關閉動作和釋放資源。

1.2.6 指定位移消費

KafkaConsumer可通過兩種方式實現實現不同粒度的指定位移消費。第一種是通過auto.offset.reset參數,另一種通過一個重要的方法seek。

(1)auto.offset.reset

auto.offset.reset這個參數總共有三種可配置的值:latest、earliest、none。如果配置不在這三個值當中,就會拋出ConfigException。

latest:當各分區下有已提交的offset時,從提交的offset開始消費;無提交的offset或位移越界時,消費新產生的該分區下的數據

earliest:當各分區下有已提交的offset時,從提交的offset開始消費;無提交的offset或位移越界時,從頭開始消費

none:topic各分區都存在已提交的offset時,從offset后開始消費;只要有一個分區不存在已提交的offset或位移越界,則拋出NoOffsetForPartitionException異常

消息的消費是通過poll方法進行的,poll方法對於開發者來說就是一個黑盒,無法精確的掌控消費的起始位置。即使通過auto.offsets.reset參數也只能在找不到位移或者位移越界的情況下粗粒度的從頭開始或者從末尾開始。因此,Kafka提供了另一種更細粒度的消費掌控:seek。

(2)seek

seek可以實現追前消費和回溯消費:

  
public void seek(TopicPartition partition, long offset)

可以通過seek方法實現指定分區的消費位移的控制。需要注意的一點是,seek方法只能重置消費者分配到的分區的偏移量,而分區的分配是在poll方法中實現的。因此在執行seek方法之前需要先執行一次poll方法獲取消費者分配到的分區,但是並不是每次poll方法都能獲得數據,所以可以採用如下的方法。

  
consumer.subscribe(topicList);
  Set<TopicPartition> assignment = new HashSet<>();
  while(assignment.size() == 0) {
      consumer.poll(Duration.ofMillis(100));
      assignment = consumer.assignment();//獲取消費者分配到的分區,沒有獲取返回一個空集合
  }

  for (TopicPartition tp : assignment) {
      consumer.seek(tp, 10); //重置指定分區的位移
  }
  while (true) {
      ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
      //consume record
    }

如果對未分配到的分區執行了seek方法,那麼會報出IllegalStateException異常。

在前面我們已經提到,使用auto.offsets.reset參數時,只有當消費者分配到的分區沒有提交的位移或者位移越界時,才能從earliest消費或者從latest消費。seek方法可以彌補這一中情況,實現任意情況的從頭或從尾部消費。

   Set<TopicPartition> assignment = new HashSet<>();
  while(assignment.size() == 0) {
      consumer.poll(Duration.ofMillis(100));
      assignment = consumer.assignment();
  }
  Map<TopicPartition, Long> offsets = consumer.endOffsets(assignment);//獲取指定分區的末尾位置
  for (TopicPartition tp : assignment) {
      consumer.seek;
  }

與endOffset對應的方法是beginningOffset方法,可以獲取指定分區的起始位置。其實kafka已經提供了一個從頭和從尾消費的方法。

  
public void seekToBeginning(Collection<TopicPartition> partitions)
public void seekToEnd(Collection<TopicPartition> partitions)

還有一種場景是這樣的,我們並不知道特定的消費位置,卻知道一個相關的時間點。為解決這種場景遇到的問題,kafka提供了一個offsetsForTimes()方法,通過時間戳來查詢分區消費的位移。

      Map<TopicPartition, Long> timestampToSearch = new HashMap<>();
  for (TopicPartition tp : assignment) {
      timestampToSearch.put(tp, System.currentTimeMillis() - 24 * 3600 * 1000);
  }
//獲得指定分區指定時間點的消費位移
  Map<TopicPartition, OffsetAndTimestamp> offsets =                                                                                   consumer.offsetsForTimes(timestampToSearch);
  for (TopicPartition tp : assignment) {
      OffsetAndTimestamp offsetAndTimestamp = offsets.get(tp);
      if (offsetAndTimestamp != null) {
              consumer.seek(tp, offsetAndTimestamp.offset());
      }
  }

由於seek方法的存在,使得消費者的消費位移可以存儲在任意的存儲介質中,包括DB、文件系統等。

1.2.7 消費者的再均衡

再均衡是指分區的所屬權從一個消費者轉移到另一消費者的行為,它為消費者組具備高可用伸縮性提高保障。不過需要注意的地方有兩點,第一是消費者發生再均衡期間,消費者組中的消費者是無法讀取消息的。第二點就是消費者發生再均衡可能會引起重複消費問題,所以一般情況下要盡量避免不必要的再均衡。

KafkaConsumer的subscribe方法中有一個參數為ConsumerRebalanceListener,我們稱之為再均衡監聽器,它可以用來在設置發生再均衡動作前後的一些準備和收尾動作。

  public interface ConsumerRebalanceListener {
  void onPartitionsRevoked(Collection<TopicPartition> partitions);
  void onPartitionsAssigned(Collection<TopicPartition> partitions);
}

onPartitionsRevoked方法會在再均衡之前和消費者停止讀取消息之後被調用。可以通過這個回調函數來處理消費位移的提交,以避免重複消費。參數partitions表示再均衡前分配到的分區。

onPartitionsAssigned方法會在再均衡之後和消費者消費之間進行調用。參數partitons表示再均衡之後所分配到的分區。

  consumer.subscribe(topicList);
  Map<TopicPartition, OffsetAndMetadata> currentOffsets = new HashMap<>();
  consumer.subscribe(topicList, new ConsumerRebalanceListener() {
      @Override
      public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
          consumer.commitSync(currentOffsets);//提交偏移量
      }

      @Override
      public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
          //do something
      }
  });

  try {
      while (isRunning.get()) {
          ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
          for (ConsumerRecord<String, String> record : records) {
              //process records
              //記錄當前的偏移量
              currentOffsets.put(new TopicPartition(record.topic(), record.partition()),new                               OffsetAndMetadata( record.offset() + 1));
          }
          consumer.commitAsync(currentOffsets, null);
      }

      } catch (Exception e) {
          e.printStackTrace();
      }finally {
          consumer.close();
      }

1.2.8 消費者攔截器

消費者攔截器主要是在消費到消息或者提交消費位移時進行一些定製化的操作。消費者攔截器需要自定義實現org.apache.kafka.clients.consumer.ConsumerInterceptor接口。

  public interface ConsumerInterceptor<K, V> extends Configurable {    
  public ConsumerRecords<K, V> onConsume(ConsumerRecords<K, V> records);
  public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets);
  public void close();
}

onConsume方法是在poll()方法返回之前被調用,比如修改消息的內容、過濾消息等。如果onConsume方法發生異常,異常會被捕獲並記錄到日誌中,但是不會向上傳遞。

Kafka會在提交位移之後調用攔截器的onCommit方法,可以使用這個方法來記錄和跟蹤消費的位移信息。

  
public class ConsumerInterceptorTTL implements ConsumerInterceptor<String,String> {
  private static final long EXPIRE_INTERVAL = 10 * 1000; //10秒過期
  @Override
  public ConsumerRecords<String, String> onConsume(ConsumerRecords<String, String> records) {
      long now = System.currentTimeMillis();
      Map<TopicPartition, List<ConsumerRecord<String, String>>> newRecords = new HashMap<>();

      for (TopicPartition tp : records.partitions()) {
          List<ConsumerRecord<String, String>> tpRecords = records.records(tp);
          List<ConsumerRecord<String, String>> newTpRecords = records.records(tp);
          for (ConsumerRecord<String, String> record : tpRecords) {
              if (now - record.timestamp() < EXPIRE_INTERVAL) {//判斷是否超時
                  newTpRecords.add(record);
              }
          }
          if (!newRecords.isEmpty()) {
              newRecords.put(tp, newTpRecords);
          }


      }
      return new ConsumerRecords<>(newRecords);
  }

  @Override
  public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) {
      offsets.forEach((tp,offset) -> {
          System.out.println(tp + ":" + offset.offset());
      });
  }

  @Override
  public void close() {}

  @Override
  public void configure(Map<String, ?> configs) {}
}

使用這種TTL需要注意的是如果採用帶參數的位移提交方式,有可能提交了錯誤的位移,可能poll拉取的最大位移已經被攔截器過濾掉。

1.2.9 消費者的多線程實現

KafkaProducer是線程安全的,然而KafkaConsumer是非線程安全的。KafkaConsumer中的acquire方法用於檢測當前是否只有一個線程在操作,如果有就會拋出ConcurrentModifiedException。acuqire方法和我們通常所說的鎖是不同的,它不會阻塞線程,我們可以把它看做是一個輕量級的鎖,它通過線程操作計數標記的方式來檢測是否發生了併發操作。acquire方法和release方法成對出現,分表表示加鎖和解鎖。

  //標記當前正在操作consumer的線程
private final AtomicLong currentThread = new AtomicLong(NO_CURRENT_THREAD);
//refcount is used to allow reentrant access by the thread who has acquired currentThread,
//大概可以理解我加鎖的次數
private final AtomicInteger refcount = new AtomicInteger(0);
private void acquire() {
long threadId = Thread.currentThread().getId();
if (threadId != currentThread.get()&&!currentThread.compareAndSet(NO_CURRENT_THREAD, threadId))
      throw new ConcurrentModificationException("KafkaConsumer is not safe for multi-threaded access");
      refcount.incrementAndGet();
}

private void release() {
  if (refcount.decrementAndGet() == 0)
      currentThread.set(NO_CURRENT_THREAD);
}

kafkaConsumer中的每個共有方法在調用之前都會執行aquire方法,只有wakeup方法是個意外。

KafkaConsumer的非線程安全並不意味着消費消息的時候只能以單線程的方式執行。可以通過多種方式實現多線程消費。

(1)Kafka多線程消費第一種實現方式——–線程封鎖

所謂線程封鎖,就是為每個線程實例化一個KafkaConsumer對象。這種方式一個線程對應一個KafkaConsumer,一個線程(可就是一個consumer)可以消費一個或多個分區的消息。這種消費方式的併發度受限於分區的實際數量。當線程數量超過分分區數量時,就會出現線程限制額的情況。

  import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Arrays;
import java.util.Properties;

public class FirstMutiConsumerDemo {
  public static final String brokerList="node112:9092,node113:9092,node114:9092";
  public static final String topic = "topic-demo";
  public static final String groupId = "group.demo";

  public static Properties initConfig() {
      Properties prop = new Properties();
      prop.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
      prop.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
      prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
      prop.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
      prop.put(ConsumerConfig.CLIENT_DNS_LOOKUP_CONFIG, "consumer.client.di.demo");
      prop.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
      return prop;
  }

  public static void main(String[] args) {
      Properties prop = initConfig();
      int consumerThreadNum = 4;
      for (int i = 0; i < 4; i++) {
          new KafkaCoosumerThread(prop, topic).run();
      }
  }

  public static class KafkaCoosumerThread extends Thread {
  //每個消費者線程包含一個KakfaConsumer對象。
      private KafkaConsumer<String, String> kafkaConsumer;
      public KafkaCoosumerThread(Properties prop, String topic) {
          this.kafkaConsumer = new KafkaConsumer<String, String>(prop);
          this.kafkaConsumer.subscribe(Arrays.asList(topic));
      }

      @Override
      public void run() {
          try {
              while (true) {
                  ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(100));
                  for (ConsumerRecord<String, String> record : records) {
                      //處理消息模塊
                  }
              }
          } catch (Exception e) {
              e.printStackTrace();
          }finally {
              kafkaConsumer.close();
          }
      }
  }
}

這種實現方式和開啟多個消費進程的方式沒有本質的區別,優點是每個線程可以按照順序消費消費各個分區的消息。缺點是每個消費線程都要維護一個獨立的TCP連接,如果分區數和線程數都很多,那麼會造成不小的系統開銷。

(2)Kafka多線程消費第二種實現方式——–多個消費線程同時消費同一分區

多個線程同時消費同一分區,通過assign方法和seek方法實現。這樣就可以打破原有消費線程個數不能超過分區數的限制,進一步提高了消費的能力,但是這種方式對於位移提交和順序控制的處理就會變得非常複雜。實際生產中很少使用。

(3)第三種實現方式——-創建一個消費者,records的處理使用多線程實現

一般而言,消費者通過poll拉取數據的速度相當快,而整體消費能力的瓶頸也正式在消息處理這一塊。基於此

考慮第三種實現方式。

  import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Arrays;
import java.util.Properties;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class ThirdMutiConsumerThreadDemo {
  public static final String brokerList="node112:9092,node113:9092,node114:9092";
  public static final String topic = "topic-demo";
  public static final String groupId = "group.demo";

  public static Properties initConfig() {
      Properties prop = new Properties();
      prop.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
      prop.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
      prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
      prop.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
      prop.put(ConsumerConfig.CLIENT_DNS_LOOKUP_CONFIG, "consumer.client.di.demo");
      prop.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
      return prop;
  }

  public static void main(String[] args) {
      Properties prop = initConfig();
      KafkaConsumerThread consumerThread = new KafkaConsumerThread(prop, topic, Runtime.getRuntime().availableProcessors());
      consumerThread.start();
  }


  public static class KafkaConsumerThread extends Thread {
      private KafkaConsumer<String, String> kafkaConsumer;
      private ExecutorService executorService;
      private int threadNum;

      public KafkaConsumerThread(Properties prop, String topic, int threadNum) {
          this.kafkaConsumer = new KafkaConsumer<String, String>(prop);
          kafkaConsumer.subscribe(Arrays.asList(topic));
          this.threadNum = threadNum;
          executorService = new ThreadPoolExecutor(threadNum, threadNum, 0L, TimeUnit.MILLISECONDS,
            new ArrayBlockingQueue<>(1000), new ThreadPoolExecutor.CallerRunsPolicy());
      }

      @Override
      public void run() {
          try {
              while (true) {
                  ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(100));
                  if (!records.isEmpty()) {
                      executorService.submit(new RecordHandler(records));
                  }
              }
          } catch (Exception e) {
              e.printStackTrace();
          }finally {
              kafkaConsumer.close();
          }
      }
  }

  public static class RecordHandler implements Runnable {
      public final ConsumerRecords<String,String> records;
      public RecordHandler(ConsumerRecords<String, String> records) {
          this.records = records;
      }
       
      @Override
      public void run() {
          //處理records
      }
  }
}

KafkaConsumerThread類對應一個消費者線程,裏面通過線程池的方式調用RecordHandler處理一批批的消息。其中線程池採用的拒絕策略為CallerRunsPolicy,當阻塞隊列填滿時,由調用線程處理該任務,以防止總體的消費能力跟不上poll拉取的速度。這種方式還可以進行橫向擴展,通過創建多個KafkaConsumerThread實例來進一步提升整體的消費能力。

這種方式還可以減少TCP連接的數量,但是對於消息的順序處理就變得困難了。這種方式需要引入一個共享變量Map<TopicPartition,OffsetAndMetadata> offsets參與消費者的偏移量提交。每一個RecordHandler類在處理完消息后都將對應的消費位移保存到共享變量offsets中,KafkaConsumerThread在每一次poll()方法之後都要進讀取offsets中的內容並對其進行提交。對於offsets的讀寫要採用加鎖處理,防止出現併發問題。並且在寫入offsets的時候需要注意位移覆蓋的問題。針對這個問題,可以將RecordHandler的run方法做如下改變:

  public void run() {
          for (TopicPartition tp : records.partitions()) {
              List<ConsumerRecord<String, String>> tpRecords = this.records.records(tp);
              //處理tpRecords
              long lastConsumedOffset = tpRecords.get(tpRecords.size() - 1).offset();
              synchronized (offsets) {
                  if (offsets.containsKey(tp)) {
                      offsets.put(tp, new OffsetAndMetadata(lastConsumedOffset + 1));
                  }else {
                      long positioin = offsets.get(tp).offset();
                      if(positioin < lastConsumedOffset + 1) {
                      offsets.put(tp, new OffsetAndMetadata(lastConsumedOffset + 1));
                      }
                  }
              }
          }
      }

對應的位移提交代碼也應該在KafkaConsumerThread的run方法中進行體現

  public void run() {
  try {
      while (true) {
          ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(100));
          if (!records.isEmpty()) {
              executorService.submit(new RecordHandler(records));
              synchronized (offsets) {
                  if (!offsets.isEmpty()) {
                      kafkaConsumer.commitSync(offsets);
                      offsets.clear();
                    }
              }
          }
      }
  } catch (Exception e) {
      e.printStackTrace();
    }finally {
        kafkaConsumer.close();
      }
    }
}

其實這種方式並不完美,可能造成數據丟失。可以通過更為複雜的滑動窗口的方式進行改進。

1.2.10 消費者重要參數

  • fetch.min.bytes

    kafkaConsumer一次拉拉取請求的最小數據量。適當增加,會提高吞吐量,但會造成額外延遲。

  • fetch.max.bytes

    kafkaConsumer一次拉拉取請求的最大數據量,如果kafka一條消息的大小超過這個值,仍然是可以拉取的。

  • fetch.max.wait.ms

    一次拉取的最長等待時間,配合fetch.min.bytes使用

  • max.partiton.fetch.bytes

    每個分區里返回consumer的最大數據量。

  • max.poll.records

    一次拉取的最大消息數

  • connection.max.idle.ms

    多久之後關閉限制的連接

  • exclude.internal.topics

    這個參數用於設置kafka中的兩個內部主題能否被公開:consumer_offsets和transaction_state。如果設為true,可以使用Pattren訂閱內部主題,如果是false,則沒有這種限制。

  • receive.buffer.bytes

    socket接收緩衝區的大小

  • send.buffer.bytes

    socket發送緩衝區的大小

  • request.timeout.ms

    consumer等待請求響應的最長時間。

  • reconnect.backoff.ms

    重試連接指定主機的等待時間。

  • max.poll.interval.ms

    配置消費者等待拉取時間的最大值,如果超過這個期限,消費者組將剔除該消費者,進行再平衡。

  • auto.offset.reset

    自動偏移量重置

  • enable.auto.commit

    是否允許偏移量的自動提交

  • auto.commit.interval.ms

    自動偏移量提交的時間間隔

【精選推薦文章】

自行創業 缺乏曝光? 下一步"網站設計"幫您第一時間規劃公司的門面形象

網頁設計一頭霧水??該從何著手呢? 找到專業技術的網頁設計公司,幫您輕鬆架站!

評比前十大台北網頁設計台北網站設計公司知名案例作品心得分享

台北網頁設計公司這麼多,該如何挑選?? 網頁設計報價省錢懶人包"嚨底家"

Android-Handler消息機制實現原理

一、消息機制流程簡介

在應用啟動的時候,會執行程序的入口函數main(),main()裏面會創建一個Looper對象,然後通過這個Looper對象開啟一個死循環,這個循環的工作是,不斷的從消息隊列MessageQueue裏面取出消息即Message對象,並處理。然後看下面兩個問題:
循環拿到一個消息之後,如何處理?
是通過在Looper的循環里調用Handler的dispatchMessage()方法去處理的,而dispatchMessage()方法裏面會調用handleMessage()方法,handleMessage()就是平時使用Handler時重寫的方法,所以最終如何處理消息由使用Handler的開發者決定。
MessageQueue里的消息從哪來?
使用Handler的開發者通過調用sendMessage()方法將消息加入到MessageQueue裏面。

上面就是Android中消息機制的一個整體流程,也是 “Android中Handler,Looper,MessageQueue,Message有什麼關係?” 的答案。通過上面的流程可以發現Handler在消息機制中的地位,是作為輔助類或者工具類存在的,用來供開發者使用。

對於這個流程有兩個疑問:

  • Looper中是如何能調用到Handler的方法的?
  • Handler是如何能往MessageQueue中插入消息的?

這兩個問題會在後面給出答案,下面先來通過源碼,分析一下這個過程的具體細節:

二、消息機制的源碼分析

首先main()方法位於ActivityThread.java類裏面,這是一個隱藏類,源碼位置:frameworks/base/core/java/android/app/ActivityThread.java

public static void main(String[] args) {
    ......
    Looper.prepareMainLooper();

    ActivityThread thread = new ActivityThread();
    thread.attach(false);

    if (sMainThreadHandler == null) {
        sMainThreadHandler = thread.getHandler();
    }

    Looper.loop();

    throw new RuntimeException("Main thread loop unexpectedly exited");
}

Looper的創建可以通過Looper.prepare()來完成,上面的代碼中prepareMainLooper()是給主線程創建Looper使用的,本質也是調用的prepare()方法。創建Looper以後就可以調用Looper.loop()開啟循環了。main方法很簡單,不多說了,下面看看Looper被創建的時候做了什麼,下面是Looper的prepare()方法和變量sThreadLocal:

static final ThreadLocal<Looper> sThreadLocal = new ThreadLocal<Looper>();

private static void prepare(boolean quitAllowed) {
    if (sThreadLocal.get() != null) {
        throw new RuntimeException("Only one Looper may be created per thread");
    }
    sThreadLocal.set(new Looper(quitAllowed));
}

很簡單,new了一個Looper,並把new出來的Looper保存到ThreadLocal裏面。ThreadLocal是什麼?它是一個用來存儲數據的類,類似HashMap、ArrayList等集合類。它的特點是可以在指定的線程中存儲數據,然後取數據只能取到當前線程的數據,比如下面的代碼:

ThreadLocal<Integer> mThreadLocal = new ThreadLocal<>();
private void testMethod() {

    mThreadLocal.set(0);
    Log.d(TAG, "main  mThreadLocal=" + mThreadLocal.get());

    new Thread("Thread1") {
        @Override
        public void run() {
            mThreadLocal.set(1);
            Log.d(TAG, "Thread1  mThreadLocal=" + mThreadLocal.get());
        }
    }.start();

    new Thread("Thread2") {
        @Override
        public void run() {
            mThreadLocal.set(2);
            Log.d(TAG, "Thread1  mThreadLocal=" + mThreadLocal.get());
        }
    }.start();

    Log.d(TAG, "main  mThreadLocal=" + mThreadLocal.get());
}

輸出的log是

main  mThreadLocal=0
Thread1  mThreadLocal=1
Thread2  mThreadLocal=2
main  mThreadLocal=0

通過上面的例子可以清晰的看到ThreadLocal存取數據的特點,只能取到當前所在線程存的數據,如果所在線程沒存數據,取出來的就是null。其實這個效果可以通過HashMap<Thread, Object>來實現,考慮線程安全的話使用ConcurrentMap<Thread, Object>,不過使用Map會有一些麻煩的事要處理,比如當一個線程結束的時候我們如何刪除這個線程的對象副本呢?如果使用ThreadLocal就不用有這個擔心了,ThreadLocal保證每個線程都保持對其線程局部變量副本的隱式引用,只要線程是活動的並且 ThreadLocal 實例是可訪問的;在線程消失之後,其線程局部實例的所有副本都會被垃圾回收(除非存在對這些副本的其他引用)。更多ThreadLocal的講解參考:Android線程管理之ThreadLocal理解及應用場景

好了回到正題,prepare()創建Looper的時候同時把創建的Looper存儲到了ThreadLocal中,通過對ThreadLocal的介紹,獲取Looper對象就很簡單了,sThreadLocal.get()即可,源碼提供了一個public的靜態方法可以在主線程的任何地方獲取這個主線程的Looper(注意一下方法名myLooper(),多個地方會用到):

public static @Nullable Looper myLooper() {
    return sThreadLocal.get();
}

Looper創建完了,接下來開啟循環,loop方法的關鍵代碼如下:

public static void loop() {
    final Looper me = myLooper();
    if (me == null) {
        throw new RuntimeException("No Looper; Looper.prepare() wasn't called on this thread.");
    }
    final MessageQueue queue = me.mQueue;

    for (;;) {
        Message msg = queue.next(); // might block
        if (msg == null) {
            // No message indicates that the message queue is quitting.
            return;
        }

        try {
            msg.target.dispatchMessage(msg);
        } finally {
            if (traceTag != 0) {
                Trace.traceEnd(traceTag);
            }
        }

        msg.recycleUnchecked();
    }
}

上面的代碼,首先獲取主線程的Looper對象,然後取得Looper中的消息隊列final MessageQueue queue = me.mQueue;,然後下面是一個死循環,不斷的從消息隊列里取消息Message msg = queue.next();,可以看到取出的消息是一個Message對象,如果消息隊列里沒有消息,就會阻塞在這行代碼,等到有消息來的時候會被喚醒。取到消息以後,通過msg.target.dispatchMessage(msg);來處理消息,msg.target 是一個Handler對象,所以這個時候就調用到我們重寫的Hander的handleMessage()方法了。
msg.target 是在什麼時候被賦值的呢?要找到這個答案很容易,msg.target是被封裝在消息裏面的,肯定要從發送消息那裡開始找,看看Message是如何封裝的。那麼就從Handler的sendMessage(msg)方法開始,過程如下:

public final boolean sendMessage(Message msg) {
    return sendMessageDelayed(msg, 0);
}

public final boolean sendMessageDelayed(Message msg, long delayMillis) {
    if (delayMillis < 0) {
        delayMillis = 0;
    }
    return sendMessageAtTime(msg, SystemClock.uptimeMillis() + delayMillis);
}

public boolean sendMessageAtTime(Message msg, long uptimeMillis) {
    MessageQueue queue = mQueue;
    if (queue == null) {
        RuntimeException e = new RuntimeException(
                this + " sendMessageAtTime() called with no mQueue");
        Log.w("Looper", e.getMessage(), e);
        return false;
    }
    return enqueueMessage(queue, msg, uptimeMillis);
}

private boolean enqueueMessage(MessageQueue queue, Message msg, long uptimeMillis) {
    msg.target = this;
    if (mAsynchronous) {
        msg.setAsynchronous(true);
    }
    return queue.enqueueMessage(msg, uptimeMillis);
}

可以看到最後的enqueueMessage()方法中msg.target = this;,這裏就把發送消息的handler封裝到了消息中。同時可以看到,發送消息其實就是往MessageQueue裏面插入了一條消息,然後Looper裏面的循環就可以處理消息了。Handler裏面的消息隊列是怎麼來的呢?從上面的代碼可以看到enqueueMessage()裏面的queue是從sendMessageAtTime傳來的,也就是mQueue。然後看mQueue是在哪初始化的,看Handler的構造方法如下:

public Handler(Callback callback, boolean async) {
    if (FIND_POTENTIAL_LEAKS) {
        final Class<? extends Handler> klass = getClass();
        if ((klass.isAnonymousClass() || klass.isMemberClass() || klass.isLocalClass()) &&
                (klass.getModifiers() & Modifier.STATIC) == 0) {
            Log.w(TAG, "The following Handler class should be static or leaks might occur: " +
                klass.getCanonicalName());
        }
    }

    mLooper = Looper.myLooper();
    if (mLooper == null) {
        throw new RuntimeException(
            "Can't create handler inside thread that has not called Looper.prepare()");
    }
    mQueue = mLooper.mQueue;
    mCallback = callback;
    mAsynchronous = async;
}

mQueue的初始化很簡單,首先取得Handler所在線程的Looper,然後取出Looper中的mQueue。這也是Handler為什麼必須在有Looper的線程中才能使用的原因,拿到mQueue就可以很容易的往Looper的消息隊列里插入消息了(配合Looper的循環+阻塞就實現了發送接收消息的效果)。

以上就是主線程中消息機制的原理。

那麼,在任何線程下使用handler的如下做法的原因、原理、內部流程等就非常清晰了:

new Thread() {
    @Override
    public void run() {
        Looper.prepare();
        Handler handler = new Handler();
        Looper.loop();
    }
}.start();
  1. 首先Looper.prepare()創建Looper並初始化Looper持有的消息隊列MessageQueue,創建好后將Looper保存到ThreadLocal中方便Handler直接獲取。
  2. 然後Looper.loop()開啟循環,從MessageQueue裏面取消息並調用handler的 dispatchMessage(msg) 方法處理消息。如果MessageQueue里沒有消息,循環就會阻塞進入休眠狀態,等有消息的時候被喚醒處理消息。
  3. 再然後我們new Handler()的時候,Handler構造方法中獲取Looper並且拿到Looper的MessageQueue對象。然後Handler內部就可以直接往MessageQueue裏面插入消息了,插入消息即發送消息,這時候有消息了就會喚醒Looper循環去處理消息。處理消息就是調用dispatchMessage(msg) 方法,最終調用到我們重寫的Handler的handleMessage()方法。

三、通過一些問題的研究加強對消息機制的理解

源碼分析完了,下面看一下文章開頭的兩個問題:

  • Looper中是如何能調用到Handler的方法的?
  • Handler是如何能往MessageQueue中插入消息的?

這兩個問題源碼分析中已經給出答案,這裏做一下總結,首先搞清楚以下對象在消息機制中的關係:

Looper,MessageQueue,Message,ThreadLocal,Handler
  1. Looper對象有一個成員MessageQueue,MessageQueue是一個消息隊列,用來存儲消息Message
  2. Message消息中帶有一個handler對象,所以Looper取出消息后,可以很方便的調用到Handler的方法(問題1解決)
  3. Message是如何帶有handler對象的?是handler在發送消息的時候把自己封裝到消息里的。
  4. Handler是如何發送消息的?是通過獲取Looper對象從而取得Looper裏面的MessageQueue,然後Handler就可以直接往MessageQueue裏面插入消息了。(問題2解決)
  5. Handler是如何獲取Looper對象的?Looper在創建的時候同時把自己保存到ThreadLocal中,並提供一個public的靜態方法可以從ThreadLocal中取出Looper,所以Handler的構造方法里可以直接調用靜態方法取得Looper對象。

帶着上面的一系列問題看源碼就很清晰了,下面是知乎上的一個問答:

Android中為什麼主線程不會因為Looper.loop()里的死循環卡死?

原因很簡單,循環里有阻塞,所以死循環並不會一直執行,相反的,大部分時間是沒有消息的,所以主線程大多數時候都是處於休眠狀態,也就不會消耗太多的CPU資源導致卡死。

  1. 阻塞的原理是使用Linux的管道機制實現的
  2. 主線程沒有消息處理時阻塞在管道的讀端
  3. binder線程會往主線程消息隊列里添加消息,然後往管道寫端寫一個字節,這樣就能喚醒主線程從管道讀端返回,也就是說looper循環里queue.next()會調用返回…

這裏說到binder線程,具體的實現細節不必深究,考慮下面的問題:
主線程的死循環如何處理其它事務?
首先需要看懂這個問題,主線程進入Looper死循環后,如何處理其他事務,比如activity的各個生命周期的回調函數是如何被執行到的(注意這裡是在同一個線程下,代碼是按順序執行的,如果在死循環這阻塞了,那麼進入死循環后循環以外的代碼是如何執行的)。
首先再看main函數的源碼

Looper.prepareMainLooper();

ActivityThread thread = new ActivityThread();
thread.attach(false);

if (sMainThreadHandler == null) {
    sMainThreadHandler = thread.getHandler();
}

Looper.loop();

在Looper.prepare和Looper.loop之間new了一個ActivityThread並調用了它的attach方法,這個方法就是開啟binder線程的,另外new ActivityThread()的時候同時會初始化它的一個H類型的成員,H是一個繼承了Handler的類。此時的結果就是:在主線程開啟loop死循環之前,已經啟動binder線程,並且準備好了一個名為H的Handler,那麼接下來在主線程死循環之外做一些事務處理就很簡單了,只需要通過binder線程向H發送消息即可,比如發送 H.LAUNCH_ACTIVITY 消息就是通知主線程調用Activity.onCreate() ,當然不是直接調用,H收到消息後會進行一系列複雜的函數調用最終調用到Activity.onCreate()。
至於誰來控制binder線程來向H發消息就不深入研究了,下面是《Android開發藝術探索》裏面的一段話:

ActivityThread 通過 ApplicationThread 和 AMS 進行進程間通訊,AMS 以進程間通信的方式完成 ActivityThread 的請求後會回調 ApplicationThread 中的 Binder 方法,然後 ApplicationThread 會向 H 發送消息,H 收到消息後會將 ApplicationThread 中的邏輯切換到 ActivityThread 中去執行,即切換到主線程中去執行,這個過程就是主線程的消息循環模型。

這個問題就到這裏,更多內容看知乎原文

最後

和其他系統相同,Android應用程序也是依靠消息驅動來工作的。網上的這句話還是很有道理的。

文章參考:

《Android開發藝術探索》
Android中為什麼主線程不會因為Looper.loop()里的死循環卡死?
Android線程管理之ThreadLocal理解及應用場景
Android 消息機制——你真的了解Handler
Android Handler到底是什麼

【精選推薦文章】

智慧手機時代的來臨,RWD網頁設計已成為網頁設計推薦首選

想知道網站建置、網站改版該如何進行嗎?將由專業工程師為您規劃客製化網頁設計及後台網頁設計

帶您來看台北網站建置台北網頁設計,各種案例分享

廣告預算用在刀口上,網站設計公司幫您達到更多曝光效益

C#中多線程中變量研究

今天在知乎上看到一個問題【為什麼在同一進程中創建不同線程,但線程各自的變量無法在線程間互相訪問?】。在多線程中,每個線程都是獨立運行的,不同的線程有可能是同一段代碼,但不會是同一作用域,所以不會共享。而共享內存,並沒有作用域之分,同一進程內,不管什麼線程都可以通過同一虛擬內存地址來訪問,不同進程也可以通過ipc等方式共享內存數據。全局變量:任何線程都可以訪問;局部變量(棧變量):任何線程執行到該函數時均可訪問,函數外不可訪問;線程變量:每個線程只能訪問自己的那個拷貝,其他線程不可見。今天就用C#來實現同一段代碼的不同線程,全局變量、局部變量、線程變量。

了解進程與線程

什麼是多任務,簡單來說就是操作系統同時可以運行多個任務。例如:一遍聽歌,一遍寫文檔等。多核CPU可以執行多任務,但是單核CPU也可以執行多任務,CPU是順序執行的,操作系統讓任務輪流執行,例如:聽歌執行一次,停頓0.01s,寫文檔執行一次,停頓0.01s等等。由於CPU的執行速度很快,我們感覺就像所有的任務都是同時執行。對操作系統來說,一個任務就是一個進程,一個進程至少有一個線程。進程是資源分配的最小單位,線程是CPU調度的最小單位。

普通的程序寫法

private static List<int> data = Enumerable.Range(1, 1000).ToList();

public static void SimpleTest()
{
    for (int i = 0; i < 10; i++)
    {
        List<int> tempData = new List<int>();
        foreach (var d in data)
        {
            tempData.Add(d);
        }
        Console.WriteLine($"i:{i},合計:{data.Sum()},是否相等:{data.Sum() == tempData.Sum()}");
    }

    Console.WriteLine("單線程運行結束");
}

多線程寫法

private static List<int> data = Enumerable.Range(1, 1000).ToList();

public static async Task MoreTaskTestAsync()
{
    List<Task> tasks = new List<Task>();
    for (int i = 0; i < 10; i++)
    {
        var tempi = i;
        var t = Task.Run(() =>
        {
            List<int> tempData = new List<int>();
            foreach (var d in data)
            {
                tempData.Add(d);
            }
            Console.WriteLine($"i:{tempi},合計:{data.Sum()},是否相等:{data.Sum() == tempData.Sum()}");
        });
        tasks.Add(t);
    }

    await Task.WhenAll(tasks); //或者Task.WaitAll(tasks.ToArray());
    Console.WriteLine("多線程運行結束");
}

不同的線程同一段代碼,但不會是同一作用域,所以tempData數據沒有互相影響。

全局變量:data,多個線程都可以訪問,list只讀的時候是線性安全
局部變量:i就是局部變量,訪問的線程可以訪問,去掉【var tempi = i;】,運行結果打印出來,值都是一樣的,增加的都是每個線程都訪問單獨的tempi變量

i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True
i:10,合計:500500,是否相等:True

線程變量:tempData,每個線程只訪問自己的,互不影響,運行結果

i:3,合計:500500,是否相等:True
i:6,合計:500500,是否相等:True
i:0,合計:500500,是否相等:True
i:1,合計:500500,是否相等:True
i:4,合計:500500,是否相等:True
i:2,合計:500500,是否相等:True
i:7,合計:500500,是否相等:True
i:5,合計:500500,是否相等:True
i:8,合計:500500,是否相等:True
i:9,合計:500500,是否相等:True

寫多線程的時候需要注意,變量的作用域,否則程序運行出來的結果將不會是想要的結果,注意,注意變量作用域。

【精選推薦文章】

自行創業 缺乏曝光? 下一步"網站設計"幫您第一時間規劃公司的門面形象

網頁設計一頭霧水??該從何著手呢? 找到專業技術的網頁設計公司,幫您輕鬆架站!

評比前十大台北網頁設計台北網站設計公司知名案例作品心得分享

台北網頁設計公司這麼多,該如何挑選?? 網頁設計報價省錢懶人包"嚨底家"

MySQL數據庫詳解(三)MySQL的事務隔離剖析

提到事務,你肯定不陌生,和數據庫打交道的時候,我們總是會用到事務。最經典的例子就是轉賬,你要給朋友小王轉 100 塊錢,而此時你的銀行卡只有 100 塊錢。

轉賬過程具體到程序里會有一系列的操作,比如查詢餘額、做加減法、更新餘額等,這些操作必須保證是一體的,不然等程序查完之後,還沒做減法之前,你這 100 塊錢,完全可以藉著這個時間差再查一次,然後再給另外一個朋友轉賬,如果銀行這麼整,不就亂了么?這時就要用到“事務”這個概念了。

簡單來說,事務就是要保證一組數據庫操作,要麼全部成功,要麼全部失敗。在 MySQL 中,事務支持是在引擎層實現的。你現在知道,MySQL 是一個支持多引擎的系統,但並不是所有的引擎都支持事務。比如 MySQL 原生的 MyISAM 引擎就不支持事務,這也是 MyISAM 被InnoDB 取代的重要原因之一。

今天的文章里,我將會以 InnoDB 為例,剖析 MySQL 在事務支持方面的特定實現,並基於原理給出相應的實踐建議,希望這些案例能加深你對 MySQL 事務原理的理解。

隔離性與隔離級別

提到事務,你肯定會想到 ACID(Atomicity、Consistency、Isolation、Durability,即原子性、一致性、隔離性、持久性),今天我們就來說說其中 I,也就是“隔離性”。

當數據庫上有多個事務同時執行的時候,就可能出現臟讀(dirty read)、不可重複讀(non reapeatable read)、幻讀(phantom read)的問題,為了解決這些問題,就有了“隔離級別”的概念。

在談隔離級別之前,你首先要知道,你隔離得越嚴實,效率就會越低。因此很多時候,我們都要在二者之間尋找一個平衡點。SQL 標準的事務隔離級別包括:讀未提交(read uncommitted)、讀提交(read committed)、可重複讀(repeatable read)和串行化serializable )。

下面我逐一為你解釋:

  • 讀未提交是指,一個事務還沒提交時,它做的變更就能被別的事務看到。
  • 讀提交是指,一個事務提交之後,它做的變更才會被其他事務看到。
  • 可重複讀是指,一個事務執行過程中看到的數據,總是跟這個事務在啟動時看到的數據是一致的。當然在可重複讀隔離級別下,未提交變更對其他事務也是不可見的。
  • 串行化,顧名思義是對於同一行記錄,“寫”會加“寫鎖”,“讀”會加“讀鎖”。當出現讀寫鎖衝突的時候,后訪問的事務必須等前一個事務執行完成,才能繼續執行。
  • 其中“讀提交”和“可重複讀”比較難理解,所以我用一個例子說明這幾種隔離級別。假設數據表 T 中只有一列,其中一行的值為 1,下面是按照時間順序執行兩個事務的行為。
mysql> create table T(c int) engine=InnoDB;
insert into T(c) values(1);

我們來看看在不同的隔離級別下,事務 A 會有哪些不同的返回結果,也就是圖裡面 V1、V2、V3 的返回值分別是什麼。

若隔離級別是“讀未提交”, 則 V1 的值就是 2。這時候事務 B 雖然還沒有提交,但是結果已經被 A 看到了。因此,V2、V3 也都是 2。

若隔離級別是“讀提交”,則 V1 是 1,V2 的值是 2。事務 B 的更新在提交后才能被 A 看到。所以, V3 的值也是 2。

若隔離級別是“可重複讀”,則 V1、V2 是 1,V3 是 2。之所以 V2 還是 1,遵循的就是這個要求:事務在執行期間看到的數據前後必須是一致的。

若隔離級別是“串行化”,則在事務 B 執行“將 1 改成 2”的時候,會被鎖住。直到事務 A提交后,事務 B 才可以繼續執行。所以從 A 的角度看, V1、V2 值是 1,V3 的值是 2。

在實現上,數據庫裏面會創建一個視圖,訪問的時候以視圖的邏輯結果為準。在“可重複讀”隔離級別下,這個視圖是在事務啟動時創建的,整個事務存在期間都用這個視圖。在“讀提交”隔離級別下,這個視圖是在每個 SQL 語句開始執行的時候創建的。這裏需要注意的是,“讀未提交”隔離級別下直接返回記錄上的最新值,沒有視圖概念;而“串行化”隔離級別下直接用加鎖的方式來避免并行訪問。

我們可以看到在不同的隔離級別下,數據庫行為是有所不同的。Oracle 數據庫的默認隔離級別其實就是“讀提交”,因此對於一些從 Oracle 遷移到 MySQL 的應用,為保證數據庫隔離級別的一致,你一定要記得將 MySQL 的隔離級別設置為“讀提交”。

配置的方式是,將啟動參數 transaction-isolation 的值設置成 READ-COMMITTED。你可以用show variables 來查看當前的值。

mysql> show variables like 'transaction_isolation';
+-----------------------+----------------+
| Variable_name | Value |
+-----------------------+----------------+
| transaction_isolation | READ-COMMITTED |
+-----------------------+----------------+

總結來說,存在即合理,哪個隔離級別都有它自己的使用場景,你要根據自己的業務情況來定。我想你可能會問那什麼時候需要“可重複讀”的場景呢?我們來看一個數據校對邏輯的案例。

假設你在管理一個個人銀行賬戶表。一個表存了每個月月底的餘額,一個表存了賬單明細。這時候你要做數據校對,也就是判斷上個月的餘額和當前餘額的差額,是否與本月的賬單明細一致。

你一定希望在校對過程中,即使有用戶發生了一筆新的交易,也不影響你的校對結果。這時候使用“可重複讀”隔離級別就很方便。事務啟動時的視圖可以認為是靜態的,不受其他事務更新的影響。

事務隔離的實現

理解了事務的隔離級別,我們再來看看事務隔離具體是怎麼實現的。這裏我們展開說明“可重複讀”。

在 MySQL 中,實際上每條記錄在更新的時候都會同時記錄一條回滾操作。記錄上的最新值,通過回滾操作,都可以得到前一個狀態的值。

假設一個值從 1 被按順序改成了 2、3、4,在回滾日誌裏面就會有類似下面的記錄。

當前值是 4,但是在查詢這條記錄的時候,不同時刻啟動的事務會有不同的 read-view。如圖中看到的,在視圖 A、B、C 裏面,這一個記錄的值分別是 1、2、4,同一條記錄在系統中可以存在多個版本,就是數據庫的多版本併發控制(MVCC)。對於 read-view A,要得到 1,就必須將當前值依次執行圖中所有的回滾操作得到。

同時你會發現,即使現在有另外一個事務正在將 4 改成 5,這個事務跟 read-view A、B、C 對應的事務是不會衝突的。

你一定會問,回滾日誌總不能一直保留吧,什麼時候刪除呢?答案是,在不需要的時候才刪除。也就是說,系統會判斷,當沒有事務再需要用到這些回滾日誌時,回滾日誌會被刪除。

什麼時候才不需要了呢?就是當系統里沒有比這個回滾日誌更早的 read-view 的時候。基於上面的說明,我們來討論一下為什麼建議你盡量不要使用長事務。

長事務意味着系統裏面會存在很老的事務視圖。由於這些事務隨時可能訪問數據庫裏面的任何數據,所以這個事務提交之前,數據庫裏面它可能用到的回滾記錄都必須保留,這就會導致大量佔用存儲空間。

在 MySQL 5.5 及以前的版本,回滾日誌是跟數據字典一起放在 ibdata 文件里的,即使長事務最終提交,回滾段被清理,文件也不會變小。我見過數據只有 20GB,而回滾段有 200GB 的庫。最終只好為了清理回滾段,重建整個庫。

除了對回滾段的影響,長事務還佔用鎖資源,也可能拖垮整個庫,這個我們會在後面講鎖的時候展開。

事務的啟動方式

如前面所述,長事務有這些潛在風險,我當然是建議你盡量避免。其實很多時候業務開發同學並不是有意使用長事務,通常是由於誤用所致。MySQL 的事務啟動方式有以下幾種:

  1. 顯式啟動事務語句, begin 或 start transaction。配套的提交語句是 commit,回滾語句
    是 rollback。
  2. set autocommit=0,這個命令會將這個線程的自動提交關掉。意味着如果你只執行一個select 語句,這個事務就啟動了,而且並不會自動提交。這個事務持續存在直到你主動執行commit 或 rollback 語句,或者斷開連接。

有些客戶端連接框架會默認連接成功后先執行一個 set autocommit=0 的命令。這就導致接下來的查詢都在事務中,如果是長連接,就導致了意外的長事務。

因此,我會建議你總是使用 set autocommit=1, 通過顯式語句的方式來啟動事務。

但是有的開發同學會糾結“多一次交互”的問題。對於一個需要頻繁使用事務的業務,第二種方式每個事務在開始時都不需要主動執行一次 “begin”,減少了語句的交互次數。如果你也有這個顧慮,我建議你使用 commit work and chain 語法。

在 autocommit 為 1 的情況下,用 begin 顯式啟動的事務,如果執行 commit 則提交事務。如果執行 commit work and chain,則是提交事務並自動啟動下一個事務,這樣也省去了再次執行 begin 語句的開銷。同時帶來的好處是從程序開發的角度明確地知道每個語句是否處於事務中。

你可以在 information_schema 庫的 innodb_trx 這個表中查詢長事務,比如下面這個語句,用於查找持續時間超過 60s 的事務。

select * from information_schema.innodb_trx where TIME_TO_SEC(timediff(now(),trx_started))>60

小結

這篇文章裏面,我介紹了 MySQL 的事務隔離級別的現象和實現,根據實現原理分析了長事務存在的風險,以及如何用正確的方式避免長事務。希望我舉的例子能夠幫助你理解事務,並更好地使用 MySQL 的事務特性。

我給你留一個問題吧。你現在知道了系統裏面應該避免長事務,如果你是業務開發負責人同時也是數據庫負責人,你會有什麼方案來避免出現或者處理這種情況呢?

【精選推薦文章】

智慧手機時代的來臨,RWD網頁設計已成為網頁設計推薦首選

想知道網站建置、網站改版該如何進行嗎?將由專業工程師為您規劃客製化網頁設計及後台網頁設計

帶您來看台北網站建置台北網頁設計,各種案例分享

廣告預算用在刀口上,網站設計公司幫您達到更多曝光效益

.NET Core 微服務之Polly重試策略

接着上一篇說,正好也是最近項目里用到了,正好拿過來整理一下,園子里也有一些文章介紹比我詳細。

簡單介紹一下紹輕量的故障處理庫 Polly  Polly是一個.NET彈性和瞬態故障處理庫

允許我們以非常順暢和線程安全的方式來執行諸如重試、斷路器、超時、隔離、緩存、後退等策略, 能為我們在微服務架構提供更穩定的服務。當然,目前的 Service Mesh 顯得更高大上,而且更強大,它更偏向從運維層面解決以上問題,不過這還是的看具體項目中怎麼去使用和決定了。

 

在微服務架構下,我們可能會遇到類似以下問題:

  1. 某些接口異常,最終造成應用程序池奔潰;
  2. 某些接口不穩定、偶爾超時,數據獲取異常;
  3. 某些服務不穩定,調用方連接不上;
  4. 某些服務異常,最終主服務掛掉(雪崩效應);

 

 當然在實際情況下,我們可能只需要確保提供給用戶的服務是可用狀態,不出現 “Service Unavailable” 這樣的畫面就好。至於接口偶爾異常,可能對某些類型的項目來說並不太關鍵,用戶可能通過重新請求、刷新頁面就可以解決,當然我們還可以在代碼層面做兼容,滿滿的try/catch、for/while 循環解決重試來保證更高的可靠性。

 這個時候Polly就能很好的起來作用,Polly 的使用相對比較簡單,當然還是得看項目結構。我們的主項目在調用微服務接口時使用了AOP,類似這種情況下,所以調用微服務的接口都是統一入口,所以我們只需要在AOP內加上 Polly 的一些策略,其他代碼不用做任何修改,就可以解決一些問題了。

安裝

Install-Package Polly

使用步驟說明

  1. 定義策略
  2. 執行方法

可以看一下代碼,我們項目主要使用的是Grpc這個框架,其他的微服務框架,使用起來大致差不多
public void Intercept(IInvocation invocation)
{
    // some code 
    try
    {
        // 創建一個策略,如果 invocation.Proceed 的執行出現 Grpc.Core.RpcException 異常,並且 StatusCode == Grpc.Core.StatusCode.Unavailable,則重試一次
        var policy = Policy
        .Handle<Grpc.Core.RpcException>(t => t.Status.StatusCode == Grpc.Core.StatusCode.Unavailable)
        .Retry(); // 默認一次

        // 將策略應用到 invocation.Proceed 方法上
        policy.Execute(invocation.Proceed);
    }
    catch (Exception ex)
    {
        // some code 
        Console.WriteLine($"{ ex.Message},{ex.StackTrace}");
    }
}

 

 

策略條件定義

策略的執行需要依賴於條件,Polly 支持對異常與結果進行策略條件定義。

異常

// 指定某個異常
Policy
  .Handle<SomeExceptionType>();

// 指定某個異常條件
Policy
  .Handle<SomeExceptionType>(ex => ex.xxx == "xxx")

// 指定多個異常
Policy
  .Handle<SomeExceptionType1>()
  .Or<SomeExceptionType2>()

// 指定多個可能異常條件
Policy
  .Handle<SomeExceptionType1>(ex => ex.xxx1 == "xxx")
  .Or<SomeExceptionType2>(ex => ex.xxx2 == "xxx")

返回結果

// 指定某個結果
Policy
  .HandleResult<ResponseMessage>(r => r.xxx == "xxx")

// 指定多個可能的結果
Policy
  .HandleResult<ResponseMessage>(r => r.xxx1 == "xxx")
  .OrResult<ResponseMessage>(r => r.xxx2 == "xxx")

重試策略(Retry )

// 指定異常下重試一次
Policy
  .Handle<SomeExceptionType>()
  .Retry();

// 指定異常下重試3次
Policy
  .Handle<SomeExceptionType>()
  .Retry(3);

// 指定異常下無限重試
Policy
  .Handle<SomeExceptionType>()
  .RetryForever();

// 每次重試之間等待指定的時間間隔
Policy
  .Handle<SomeExceptionType>()
  .WaitAndRetry(new[]
  {
    TimeSpan.FromSeconds(1),
    TimeSpan.FromSeconds(3),
    TimeSpan.FromSeconds(7)
  });

Retry 可以指定一個要執行的 Action。Action 參數:exception 當前異常信息,retryCount 當前執行第幾次,context 當前執行上下文信息。

測試一下:

private static int times = 0;

public static void TestPolicy()
{
    var policy = Policy
        .Handle<Exception>()
        .Retry(3, (exception, retryCount, context) => // 出異常會執行以下代碼
        {
            Console.WriteLine($"exception:{ exception.Message}, retryCount:{retryCount}, id:{context["id"]}, name:{context["name"]}");
        });

    try
    {
        // 通過 new Context 傳遞上下文信息
        var result = policy.Execute(Test, new Context("data", new Dictionary<string, object>() { { "id", "1" }, { "name", "beck" } }));
        Console.WriteLine($"result:{result}");
    }
    catch (Exception ex)
    {
        Console.WriteLine(ex.Message);
    }
}

private static string Test()
{
    // 每執行一次加1
    times++;

    // 前2次都拋異常
    if (times < 3)
    {
        throw new Exception("exception message");
    }
    return "success";
}

測試結果:

 

可以看到得到了咱們想要的效果,具體項目可以具體去實施,下一篇咱們接着說Polly的熔斷策略。感興趣可以自行搜索Polly的相關文檔看看。

參考鏈接

  • Polly
  • Polly Project
  • PollySamples

 

沒有彩蛋

 

【精選推薦文章】

自行創業 缺乏曝光? 下一步"網站設計"幫您第一時間規劃公司的門面形象

網頁設計一頭霧水??該從何著手呢? 找到專業技術的網頁設計公司,幫您輕鬆架站!

評比前十大台北網頁設計台北網站設計公司知名案例作品心得分享

台北網頁設計公司這麼多,該如何挑選?? 網頁設計報價省錢懶人包"嚨底家"

java8 函數式接口——Function/Predict/Supplier/Consumer

Function

我們知道Java8的最大特性就是函數式接口。所有標註了@FunctionalInterface註解的接口都是函數式接口,具體來說,所有標註了該註解的接口都將能用在lambda表達式上。

接口介紹

/**
 * Represents a function that accepts one argument and produces a result.
 *
 * <p>This is a <a href="package-summary.html">functional interface</a>
 * whose functional method is {@link #apply(Object)}.
 *
 * @param <T> the type of the input to the function
 * @param <R> the type of the result of the function
 *
 * @since 1.8
 */

上述描述可知: Function中傳遞的兩個泛型:T,R分別代表 輸入參數類型和返回參數類型。下面將逐個介紹Function中的各個接口:

接口1: 執行具體內容接口
R apply(T t);

實例1:apply使用
    // 匿名類的方式實現
    Function<Integer, Integer> version1 = new Function<Integer, Integer>() {
        @Override
        public Integer apply(Integer integer) {
            return integer++;
        }
    };
    int result1 = version1.apply(20);


    // lamda表達式
    Function<Integer, Integer> version2 = integer -> integer++;
    int result2 = version1.apply(20);
    

接口2: compose
該方法是一個默認方法,這個方法接收一個function作為參數,將參數function執行的結果作為參數給調用的function,以此來實現兩個function組合的功能。

// compose 方法源碼
default <V> Function<V, R> compose(Function<? super V, ? extends T> before) {
        Objects.requireNonNull(before);
        return (V v) -> apply(before.apply(v));
    }
實例2:compose使用
public int compute(int a, Function<Integer, Integer> function1, Function<Integer, Integer> function2) {
    return function1.compose(function2).apply(a);
}

// 調用上述方法
test.compute(2, value -> value * 3, value -> value * value) 
// 執行結果: 12 (有源碼可以看出先執行before)

接口3 : andThen
了解了compose方法,我們再來看andThen方法就好理解了,聽名字就是“接下來”,andThen方法也是接收一個function作為參數,與compse不同的是,先執行本身的apply方法,將執行的結果作為參數給參數中的function。

public interface Function<T, R> {
    default <V> Function<T, V> andThen(Function<? super R, ? extends V> after) {
        Objects.requireNonNull(after);
        return (T t) -> after.apply(apply(t));
    }
}
實例3:andThen使用
public int compute2(int a, Function<Integer, Integer> function1, Function<Integer, Integer> function2) {
    return function1.andThen(function2).apply(a);
}

// 調用上述方法
test.compute2(2, value -> value * 3, value -> value * value) 
// 執行結果 : 36

反思: 多個參數

Function接口雖然很簡潔,但是由Function源碼可以看出,他只能傳一個參數,實際使用中肯定不能滿足需求。下面提供幾種思路:

  1. BiFunction可以傳遞兩個參數(Java8中還提供了其它相似Function)
  2. 通過封裝類來解決
  3. void函數還是無法解決

因為參數原因“自帶的Function”函數必然不能滿足業務上複雜多變的需求,那麼就自定義Function接口吧

@FunctionalInterface
    static interface ThiConsumer<T,U,W>{
        void accept(T t, U u, W w);

        default ThiConsumer<T,U,W> andThen(ThiConsumer<? super T,? super U,? super W> consumer){
            return (t, u, w)->{
                accept(t, u, w);
                consumer.accept(t, u, w);
            };
        }
    }

自此,Function接口介紹完畢。

斷言性接口:Predicate

接口介紹:

/**
 * Represents a predicate (boolean-valued function) of one argument.
 *
 * <p>This is a <a href="package-summary.html">functional interface</a>
 * whose functional method is {@link #test(Object)}.
 *
 * @param <T> the type of the input to the predicate
 *
 * @since 1.8
 */
@FunctionalInterface
public interface Predicate<T> {

Predicate是個斷言式接口其參數是<T,boolean>,也就是給一個參數T,返回boolean類型的結果。跟Function一樣,Predicate的具體實現也是根據傳入的lambda表達式來決定的。

源碼不再具體分析,主要有 test/and/or/negate方法,以及一個靜態方法isEqual,具體使用實例如下:

    private static void testPredict() {
        int[] numbers = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15};
        List<Integer> list = new ArrayList<>();
        for (int i : numbers) {
            list.add(i);
        }
        
        // 三個判斷
        Predicate<Integer> p1 = i -> i > 5;
        Predicate<Integer> p2 = i -> i < 20;
        Predicate<Integer> p3 = i -> i % 2 == 0;
        List test = list.stream()
                .filter(p1
                        .and(p2)
//                        .and(Predicate.isEqual(8))
                        .and(p3))
                .collect(Collectors.toList());
        System.out.println(test.toString());
        //print:[6, 8, 10, 12, 14]
    }

供給性接口:Supplier

接口介紹

/**
 * Represents a supplier of results.
 *
 * <p>There is no requirement that a new or distinct result be returned each
 * time the supplier is invoked.
 *
 * <p>This is a <a href="package-summary.html">functional interface</a>
 * whose functional method is {@link #get()}.
 *
 * @param <T> the type of results supplied by this supplier
 *
 * @since 1.8
 */
@FunctionalInterface
public interface Supplier<T> 

使用實例:

        Supplier supplier = "Hello"::toLowerCase;
        System.out.println(supplier);

消費性:Consumer

接口介紹

/**
 * Represents an operation that accepts a single input argument and returns no
 * result. Unlike most other functional interfaces, {@code Consumer} is expected
 * to operate via side-effects.
 *
 * <p>This is a <a href="package-summary.html">functional interface</a>
 * whose functional method is {@link #accept(Object)}.
 *
 * @param <T> the type of the input to the operation
 *
 * @since 1.8
 */
@FunctionalInterface
public interface Consumer<T> {

實際使用

    NameInfo info = new NameInfo("abc", 123);
    Consumer<NameInfo> consumer = t -> {
        String infoString = t.name + t.age;
        System.out.println("consumer process:" + infoString);
    };
    consumer.accept(info);

總結:
本文主要介紹Java8的接口式編程,以及jdk中提供的四種函數接口(FunctionalInterface)。Predict/Supplier/Consumer其實是Function的一種變形,所以沒有詳細介紹。
疑問: FunctionalInterface註解是如何和lamada表達式聯繫在一起,函數接口在編譯時又是如何處理的?後面再了解下

【精選推薦文章】

智慧手機時代的來臨,RWD網頁設計已成為網頁設計推薦首選

想知道網站建置、網站改版該如何進行嗎?將由專業工程師為您規劃客製化網頁設計及後台網頁設計

帶您來看台北網站建置台北網頁設計,各種案例分享

廣告預算用在刀口上,網站設計公司幫您達到更多曝光效益

SpringCloud-分佈式配置中心【入門介紹】

案例代碼:https://github.com/q279583842q/springcloud-e-book

一、 為什麼需要使用配置中心

1 服務配置的現狀

2 常用的配置管理解決方案的缺點

3 為什麼要使用 spring cloud config 配置中心?

4 spring cloud config配置中心,它解決了什麼問題?

二、 編寫配置中心入門案例

1.編寫配置中心的服務端

1.1 創建服務端項目

  創建一個SpringCloud項目。

1.2 修改pom文件

  我們需要添加config-server的依賴,具體如下

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>1.5.13.RELEASE</version>
    </parent>
    <groupId>com.bobo</groupId>
    <artifactId>config-server</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <dependencyManagement>
        <dependencies>
            <dependency>
                <groupId>org.springframework.cloud</groupId>
                <artifactId>spring-cloud-dependencies</artifactId>
                <version>Dalston.SR5</version>
                <type>pom</type>
                <scope>import</scope>
            </dependency>
        </dependencies>
    </dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-eureka</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-config-server</artifactId>
        </dependency>
    </dependencies>
    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>
    </build>
</project>

1.3 修改配置文件

  在此處的配置文件中我們需要關聯碼雲或者GitHub。以碼雲為例

碼雲處理

  首先我們需要在碼雲上註冊一個賬號(https://gitee.com) 然後創建一個新的項目。

配置文件處理

  在配置文件中添加如下配置

spring.application.name=config-server
server.port=9050
# 設置服務註冊中心地址,指向另一個註冊中心
eureka.client.serviceUrl.defaultZone=http://dpb:123456@eureka1:8761/eureka/,http://dpb:123456@eureka2:8761/eureka/

#Git 配置
spring.cloud.config.server.git.uri=https://gitee.com/dengpbs/config
#spring.cloud.config.server.git.username=
#spring.cloud.config.server.git.password=

創建四個配置文件

四個配置文件都有一個e-book屬性,只是值不一樣。然後將這個四個配置文件上傳到碼雲中我們新創建的倉庫

然後將項目中的四個配置文件刪除

1.4 修改啟動類

  我們需要在啟動類中添加eureka客戶端和config服務端的註解,具體如下:

@SpringBootApplication
@EnableEurekaClient
@EnableConfigServer
public class ConfigServerStart {

    public static void main(String[] args) {
        SpringApplication.run(ConfigServerStart.class, args);
    }
}

1.5 訪問測試

  啟動服務,訪問測試
http://localhost:9050/config-client/test

http://localhost:9050/config-client/default

http://localhost:9050/config-client/dev

通過訪問,我們獲取到了位於碼雲倉庫中的屬性信息。

1.6 配置文件的命名規則與訪問

  注意,上面案例中的配置文件的名稱,並不是隨便命名的,而是有一定的規則來約束的,具體如下:

2.編寫客戶端程序

2.1 創建項目

  創建一個SpringCloud項目

2.2 pom文件修改

  配置中心的客戶端使用的依賴需要注意,不是config-server了,具體如下:

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>1.5.13.RELEASE</version>
    </parent>
    <groupId>com.bobo</groupId>
    <artifactId>config-client</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <dependencyManagement>
        <dependencies>
            <dependency>
                <groupId>org.springframework.cloud</groupId>
                <artifactId>spring-cloud-dependencies</artifactId>
                <version>Dalston.SR5</version>
                <type>pom</type>
                <scope>import</scope>
            </dependency>
        </dependencies>
    </dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-eureka</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-config</artifactId>
        </dependency>
    </dependencies>
    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>
    </build>
</project>

2.3 修改配置文件

  注意在配置中心的客戶端服務中,配置文件的名稱必須是bootstrap.properties或者bootstrap.yml文件。
官方解釋:

Spring Cloud 構建於 Spring Boot 之上,在 Spring Boot 中有兩種上下文,一種是 bootstrap, 另外一種是 application, bootstrap 是應用程序的父上下文,也就是說 bootstrap 加載優先於 applicaton。bootstrap 主要用於從額外的資源來加載配置信息,還可以在本地外部配置文件中解密屬性。這兩個上下文共用一個環境,它是任何Spring應用程序的外部屬性的來源。bootstrap 裏面的屬性會優先加載,它們默認也不能被本地相同配置覆蓋。

spring.application.name=config-client
server.port=9051
#設置服務註冊中心地址,指向另一個註冊中心
eureka.client.serviceUrl.defaultZone=http://dpb:123456@eureka1:8761/eureka/,http://dpb:123456@eureka2:8761/eureka/

#默認 false,這裏設置 true,表示開啟讀取配置中心的配置
spring.cloud.config.discovery.enabled=true
#對應 eureka 中的配置中心 serviceId,默認是 configserver
spring.cloud.config.discovery.serviceId=config-server
#指定環境
spring.cloud.config.profile=dev
#git 標籤
spring.cloud.config.label=master

2.4 修改啟動類

@SpringBootApplication
@EnableEurekaClient
public class ConfigClientStart {

    public static void main(String[] args) {
        SpringApplication.run(ConfigClientStart.class, args);
    }
}

2.5 創建控制器

  在控制中我們嘗試獲取配置中心的數據,具體如下:

@RestController
public class ShowController {
    
    @Value("${e-book}")
    private String msg;
    
    @RequestMapping("/showMsg")
    public String showMsg(){
        return msg;
    }
}

2.6 啟動測試

訪問:http://localhost:9051/showMsg

搞定~

【精選推薦文章】

自行創業 缺乏曝光? 下一步"網站設計"幫您第一時間規劃公司的門面形象

網頁設計一頭霧水??該從何著手呢? 找到專業技術的網頁設計公司,幫您輕鬆架站!

評比前十大台北網頁設計台北網站設計公司知名案例作品心得分享

台北網頁設計公司這麼多,該如何挑選?? 網頁設計報價省錢懶人包"嚨底家"

python算法與數據結構-順序表(37)

 

1、順序表介紹

  順序表是最簡單的一種線性結構,邏輯上相鄰的數據在計算機內的存儲位置也是相鄰的,可以快速定位第幾個元素,中間不允許有空,所以插入、刪除時需要移動大量元素。順序表可以分配一段連續的存儲空間Maxsize,用elem記錄基地址,用length記錄實際的元素個數,即順序表的長度, 

  上圖1表示的是順序表的基本形式,數據元素本身連續存儲,每個元素所佔的存儲單元大小固定相同,元素的下標是其邏輯地址,而元素存儲的物理地址(實際內存地址)可以通過存儲區的起始地址Loc (e0)加上邏輯地址(第i個元素)與存儲單元大小(c)的乘積計算而得,即:Loc(element i) = Loc(e0) + c*i

所以、訪問指定元素時無需從頭遍歷,通過計算便可獲得對應地址,其時間複雜度為O(1)。

  如果元素的大小不統一,則須採用圖2的元素外置的形式,將實際數據元素另行存儲,而順序表中各單元位置保存對應元素的地址信息(即鏈接)。由於每個鏈接所需的存儲量相同,通過上述公式,可以計算出元素鏈接的存儲位置,而後順着鏈接找到實際存儲的數據元素。注意,圖2中的c不再是數據元素的大小,而是存儲一個鏈接地址所需的存儲量,這個量通常很小。

圖2這樣的順序表也被稱為對實際數據的索引,這是最簡單的索引結構。

2、順序表的結構 

  

  一個順序表的完整信息包括兩部分,一部分是表中的元素集合,另一部分是為實現正確操作而需記錄的信息,即有關表的整體情況的信息,這部分信息主要包括元素存儲區的容量和當前表中已有的元素個數兩項。

3、順序表的兩種基本實現方式

  1為一體式結構,存儲表信息的單元與元素存儲區以連續的方式安排在一塊存儲區里,兩部分數據的整體形成一個完整的順序表對象。一體式結構整體性強,易於管理。但是由於數據元素存儲區域是表對象的一部分,順序表創建后,元素存儲區就固定了。

  2為分離式結構,表對象里只保存與整個表有關的信息(即容量和元素個數),實際數據元素存放在另一個獨立的元素存儲區里,通過鏈接與基本表對象關聯。

 4、元素存儲區替換

  一體式結構由於順序表信息區與數據區連續存儲在一起,所以若想更換數據區,則只能整體搬遷,即整個順序表對象(指存儲順序表的結構信息的區域)改變了。分離式結構若想更換數據區,只需將表信息區中的數據區鏈接地址更新即可,而該順序表對象不變。

5、元素存儲區擴充

  採用分離式結構的順序表,若將數據區更換為存儲空間更大的區域,則可以在不改變表對象的前提下對其數據存儲區進行了擴充,所有使用這個表的地方都不必修改。只要程序的運行環境(計算機系統)還有空閑存儲,這種表結構就不會因為滿了而導致操作無法進行。人們把採用這種技術實現的順序表稱為動態順序表,因為其容量可以在使用中動態變化。

擴充的兩種策略

  • 每次擴充增加固定數目的存儲位置,如每次擴充增加10個元素位置,這種策略可稱為線性增長。

    特點:節省空間,但是擴充操作頻繁,操作次數多。

  • 每次擴充容量加倍,如每次擴充增加一倍存儲空間。

    特點:減少了擴充操作的執行次數,但可能會浪費空間資源。以空間換時間,推薦的方式。

6、順序表的增刪改查操作的Python代碼實現

# 創建順序表
class Sequence_Table():
    
    # 初始化
    def __init__(self):
        self.date = [None]*100
        self.length = 0
    
    # 判斷是否已經滿了
    def isFull(self):
        if self.length>100:
            print("該順序表已滿,無法添加元素")
            return 1
        else:
            return 0
    
    # 按下錶索引查找
    def selectByIndex(self,index):
        if index>=0 and index<=self.length-1:
            return self.date[index]
        else:
            print("你輸入的下標不對,請重新輸入\n")
            return 0
        
    # 按元素查下標
    def selectByNum(self,num):
        isContain = 0
        for i in range(0,self.length):
            if self.date[i] == num:
                isContain = 1
                print("你要查找的元素下標是%d\n"%i)
        if isContain == 0:
            print("沒有找你你要的數據")
    
    # 追加數據
    def addNum(self,num):
        if self.isFull() == 0:
            self.date[self.length] = num
            self.length += 1
            
    # 打印順序表
    def printAllNum(self):
        for i in range(self.length):
            print("a[%s]=%s"%(i,self.date[i]),end=" ")
        print("\n")
        
    # 按下標插入數據
    def insertNumByIndex(self,num,index):
        if index<0 or index>self.length:
            return 0
        self.length += 1
        for i in range(self.length-1,index,-1):
            temp = self.date[i]
            self.date[i] = self.date[i-1]
            self.date[i-1] = temp
        self.date[index] = num
        return 1
    # 按下標刪除數據
    def delectNumByIndex(self,index):
        if self.length <= 0:
            print("該順序表內沒有數據,不用刪除")
            
        for i in range(index,self.length-1):
            temp = self.date[i]
            self.date[i] = self.date[i + 1]
            self.date[i + 1] = temp
        self.date[self.length-1] = 0
        self.length -= 1

def main():
    # 創建順序表對象
    seq_t = Sequence_Table()
    
    # 插入三個元素
    seq_t.addNum(1)
    seq_t.addNum(2)
    seq_t.addNum(3)
    
    # 打印驗證
    seq_t.printAllNum()
    
    # 按照索引查找
    num = seq_t.selectByIndex(2)
    print("你要查找的數據是%d\n" % num)
    
    # 按照索引插入數據
    seq_t.insertNumByIndex(4, 1)
    seq_t.printAllNum()
    
    # 按照数字查下標
    seq_t.selectByNum(4)
    
    #刪除數據
    seq_t.delectNumByIndex(1)
    seq_t.printAllNum()
     
if __name__ == "__main__":
    main()

運行結果為:

a[0]=1 a[1]=2 a[2]=3 

你要查找的數據是3

a[0]=1 a[1]=4 a[2]=2 a[3]=3 

你要查找的元素下標是1

a[0]=1 a[1]=2 a[2]=3 

7、順序表的增刪改查操作的C語言代碼實現

#include<stdio.h>
// 1、定義順序表的儲存結構
typedef struct
{
    //用數組存儲線性表中的元素
    int data[100];
    // 順序表中的元素個數
    int length;
}Sequence_table,*p_Sequence_table;

// 2、順序表的初始化,
void initSequenceTable(p_Sequence_table T)
{
    // 判斷傳過來的表是否為空,為空直接退出
    if (T == NULL)
    {
        return;
    }
    // 設置默認長度為0
    T->length = 0;
}

// 3、求順序表的長度
int lengthOfSequenceTable(p_Sequence_table T)
{
    if (T==NULL)
    {
        return 0;
    }
    return T->length;
}

// 4、判斷順序表是否已滿
int isFull(p_Sequence_table T)
{
    if (T->length>=100)
    {
        printf("該順序表已經裝滿,無法再添加元素");
        return 1;
    }
    return 0;
}

// 5、按序號查找
int selectSequenceTableByIndex(p_Sequence_table T,int index)
{
    if (index>=0&&index<=T->length-1)
    {
        return T->data[index];
    }
    printf("你輸入的序號不對,請重新輸入\n");
    return 0;
}

// 6、按內容查找是否存在
void selectSequenceTableByNum(p_Sequence_table T,int num)
{
    int isContain = 0;
    for (int i=0; i<T->length; i++)
    {
        if (T->data[i] == num)
        {
            isContain = 1;
            printf("你要找的元素的下標是:%d\n",i);
        }
    }
    if (isContain == 0)
    {
        printf("沒有找到你要的數據\n");
    }
}

// 7、添加元素(在隊尾添加)
void addNumber(p_Sequence_table T,int num)
{
    // 順序表還沒有滿的時候
    if (isFull(T) == 0)
    {
        T->data[T->length] = num;
        T->length++;
    }
}

// 8、順序表的遍歷
void printAllNumOfSequenceTable(p_Sequence_table T)
{
    for (int i = 0; i<T->length; i++)
    {
        printf("T[%d]=%d ",i,T->data[i]);
    }
    printf("\n");
}

//9、插入操作
int insertNumByIndex(p_Sequence_table T, int num,int index)
{
    if (index<0||index>T->length)
    {
        return 0;
    }
    T->length++;
    for (int i = T->length-1; i>index; i--)
    {
        int temp = T->data[i];
        T->data[i] = T->data[i-1];
        T->data[i-1] = temp;
    }
    T->data[index] = num;
    return 1;
}

// 10、刪除元素
void delectNum(p_Sequence_table T,int index)
{
    if (T->length <= 0)
    {
        printf("該順序表中沒有數據,不用刪除");
    }
    for (int i = index;i<T->length-1; i++)
    {
        int temp = T->data[i];
        T->data[i] = T->data[i+1];
        T->data[i+1] = temp;
    }
    T->data[T->length-1] = 0;
    T->length--;
}



int main(int argc, const char * argv[]) {
    
    // 創建順序表的結構體
    Sequence_table seq_t;
    // 初始化
    initSequenceTable(&seq_t);
    // 添加數據
    addNumber(&seq_t, 1);
    addNumber(&seq_t, 2);
    addNumber(&seq_t, 3);
    // 打印驗證
    printAllNumOfSequenceTable(&seq_t);
    // 根據索引下標查內容
    int num = selectSequenceTableByIndex(&seq_t, 2);
    printf("你查的數據是:%d\n",num);
    // 插入
    insertNumByIndex(&seq_t, 4, 1);
    printAllNumOfSequenceTable(&seq_t);
    // 根據內容查下標
    selectSequenceTableByNum(&seq_t, 4);
    // 根據下標刪除數據
    delectNum(&seq_t, 1);
    printAllNumOfSequenceTable(&seq_t);
    return 0;
}

運行結果為:

T[0]=1 T[1]=2 T[2]=3 
你查的數據是:3
T[0]=1 T[1]=4 T[2]=2 T[3]=3 
你要找的元素的下標是:1
T[0]=1 T[1]=2 T[2]=3 

 

【精選推薦文章】

智慧手機時代的來臨,RWD網頁設計已成為網頁設計推薦首選

想知道網站建置、網站改版該如何進行嗎?將由專業工程師為您規劃客製化網頁設計及後台網頁設計

帶您來看台北網站建置台北網頁設計,各種案例分享

廣告預算用在刀口上,網站設計公司幫您達到更多曝光效益

SpringBoot之ApplicationContextInitializer的理解和使用

一、 ApplicationContextInitializer 介紹

  首先看spring官網的介紹:

   翻譯一下:

  • 用於在spring容器刷新之前初始化Spring ConfigurableApplicationContext的回調接口。(剪短說就是在容器刷新之前調用該類的 initialize 方法。並將 ConfigurableApplicationContext 類的實例傳遞給該方法)
  • 通常用於需要對應用程序上下文進行編程初始化的web應用程序中。例如,根據上下文環境註冊屬性源或激活配置文件等。
  • 可排序的(實現Ordered接口,或者添加@Order註解)

  看完這段解釋,為了講解方便,我們先看自定義 ApplicationContextInitializer 的三種方式。再通過SpringBoot的源碼,分析生效的時間以及實現的功能等。

二、三種實現方式

  首先新建一個類 MyApplicationContextInitializer 並實現 ApplicationContextInitializer 接口。

1 public class MyApplicationContextInitializer implements ApplicationContextInitializer {
2     @Override
3     public void initialize(ConfigurableApplicationContext applicationContext) {
4         System.out.println("-----MyApplicationContextInitializer initialize-----");
5     }
6 }

  2.1、mian函數中添加

  優雅的寫一個SpringBoot的main方法

1 @SpringBootApplication
2 public class MySpringBootApplication {
3     public static void main(String[] args) {
4         SpringApplication application = new SpringApplication(MySpringBootApplication.class);
5         application.addInitializers(new MyApplicationContextInitializer());
6         application.run(args);
7     }
8 }

 

  運行,查看控制台:生效了

  

  2.2、配置文件中配置

context.initializer.classes=org.springframework.boot.demo.common.MyApplicationContextInitializer 

 

  

  2.3、SpringBoot的SPI擴展—META-INF/spring.factories中配置

org.springframework.context.ApplicationContextInitializer=org.springframework.boot.demo.common.MyApplicationContextInitializer

 

  

 

三、排序問題

  如圖所示改造一下mian方法。打一個斷點,debug查看排序情況。

  

  給 MyApplicationContextInitializer 加上Order註解:我們指定其擁有最高的排序級別。(越高越早執行)

1 @Order(Ordered.HIGHEST_PRECEDENCE)
2 public class MyApplicationContextInitializer implements ApplicationContextInitializer{
3     @Override
4     public void initialize(ConfigurableApplicationContext applicationContext) {
5         System.out.println("-----MyApplicationContextInitializer initialize-----");
6     }
7 }

 

  下面我們通過debug分別驗證二章節中提到的三種方法排序是否都是可以的。

  首先驗證2.1章節中採用的main函數中添加:debug,斷點處查看 application.getInitializers() 這行代碼的結果可見,排序生效了。

  

  然後再分別驗證2.2和2.3章節中的方法。排序都是可以實現的。

  然而當採用2.3中的SPI擴展的方式,排序指定 @Order(Ordered.LOWEST_PRECEDENCE) 排序並沒有生效。當然採用實現Ordered接口的方式,排序驗證結果都是一樣的。

 四、通過源碼分析ApplicationContextInitializer何時被調用

  debug差看上文中自定的 MyApplicationContextInitializer 的調用棧。

  

  可見 ApplicationContextInitializer 在容器刷新前的準備階段被調用。 refreshContext(context); 

  在SpringBoot的啟動函數中, ApplicationContextInitializer 

 1     public ConfigurableApplicationContext run(String... args) {
 2         //記錄程序運行時間
 3         StopWatch stopWatch = new StopWatch();
 4         stopWatch.start();
 5         // ConfigurableApplicationContext Spring 的上下文
 6         ConfigurableApplicationContext context = null;
 7         Collection<SpringBootExceptionReporter> exceptionReporters = new ArrayList<>();
 8         configureHeadlessProperty();
 9         //從META-INF/spring.factories中獲取監聽器
10         //1、獲取並啟動監聽器
11         SpringApplicationRunListeners listeners = getRunListeners(args);
12         listeners.starting();
13         try {
14             ApplicationArguments applicationArguments = new DefaultApplicationArguments(
15                     args);
16             //2、構造容器環境
17             ConfigurableEnvironment environment = prepareEnvironment(listeners, applicationArguments);
18             //處理需要忽略的Bean
19             configureIgnoreBeanInfo(environment);
20             //打印banner
21             Banner printedBanner = printBanner(environment);
22             ///3、初始化容器
23             context = createApplicationContext();
24             //實例化SpringBootExceptionReporter.class,用來支持報告關於啟動的錯誤
25             exceptionReporters = getSpringFactoriesInstances(
26                     SpringBootExceptionReporter.class,
27                     new Class[]{ConfigurableApplicationContext.class}, context);
28             //4、刷新容器前的準備階段
29             prepareContext(context, environment, listeners, applicationArguments, printedBanner);
30             //5、刷新容器
31             refreshContext(context);
32             //刷新容器后的擴展接口
33             afterRefresh(context, applicationArguments);
34             stopWatch.stop();
35             if (this.logStartupInfo) {
36                 new StartupInfoLogger(this.mainApplicationClass)
37                         .logStarted(getApplicationLog(), stopWatch);
38             }
39             listeners.started(context);
40             callRunners(context, applicationArguments);
41         } catch (Throwable ex) {
42             handleRunFailure(context, ex, exceptionReporters, listeners);
43             throw new IllegalStateException(ex);
44         }
45 
46         try {
47             listeners.running(context);
48         } catch (Throwable ex) {
49             handleRunFailure(context, ex, exceptionReporters, null);
50             throw new IllegalStateException(ex);
51         }
52         return context;
53     }

 

   然後看在 refreshContext(context); 具體是怎麼被調用的。

1 private void prepareContext(ConfigurableApplicationContext context,
2                             ConfigurableEnvironment environment, SpringApplicationRunListeners listeners,
3                             ApplicationArguments applicationArguments, Banner printedBanner) {
4     context.setEnvironment(environment);
5     postProcessApplicationContext(context);
6     applyInitializers(context);
7     ...
8 }

 

   然後在 applyInitializers 中遍歷調用每一個被加載的 ApplicationContextInitializer 的  initialize(context);  方法,並將 ConfigurableApplicationContext 的實例傳遞給 initialize 方法。

1 protected void applyInitializers(ConfigurableApplicationContext context) {
2     for (ApplicationContextInitializer initializer : getInitializers()) {
3         Class<?> requiredType = GenericTypeResolver.resolveTypeArgument(
4                 initializer.getClass(), ApplicationContextInitializer.class);
5         Assert.isInstanceOf(requiredType, context, "Unable to call initializer.");
6         initializer.initialize(context);
7     }
8 }

 

  OK,到這裏通過源碼說明了 ApplicationContextInitializer 是何時及如何被調用的。

 

【精選推薦文章】

自行創業 缺乏曝光? 下一步"網站設計"幫您第一時間規劃公司的門面形象

網頁設計一頭霧水??該從何著手呢? 找到專業技術的網頁設計公司,幫您輕鬆架站!

評比前十大台北網頁設計台北網站設計公司知名案例作品心得分享

台北網頁設計公司這麼多,該如何挑選?? 網頁設計報價省錢懶人包"嚨底家"

從0到1:全面理解RPC遠程調用

上一篇關於 WSGI 的硬核長文,不知道有多少同學,能夠從頭看到尾的,不管你們有沒有看得很過癮,反正我是寫得很爽,總有一種將一樣知識吃透了的錯覺。

今天我又給自己挖坑了,打算將 rpc 遠程調用的知識,好好地梳理一下,花了周末整整两天的時間。

什麼是RPC呢?

百度百科給出的解釋是這樣的:“RPC(Remote Procedure Call Protocol)——遠程過程調用協議,它是一種通過網絡從遠程計算機程序上請求服務,而不需要了解底層網絡技術的協議”。這個概念聽起來還是比較抽象,沒關係,繼續往後看,後面概念性的東西,我會講得足夠清楚,讓你完全掌握 RPC 的基礎內容。在後面的篇章中還會結合其在 OpenStack 中實際應用,一步一步揭開 rpc 的神秘面紗。

有的讀者,可能會問,為啥我舉的例子老是 OpenStack 里的東西呢?

因為每個人的業務中接觸的框架都不一樣(我主要接觸的就是 OpenStack 框架),我無法為每個人去定製寫一篇文章,但其技術原理都是一樣的。即使如此,我也會儘力將文章寫得通用,不會因為你沒接觸過 OpenStack 而成為你理解 rpc 的瓶頸。

01. 既 REST,何 RPC ?

在 OpenStack 里的進程間通信方式主要有兩種,一種是基於HTTP協議的RESTFul API方式,另一種則是RPC調用。

那麼這兩種方式在應用場景上有何區別呢?

有使用經驗的人,就會知道:

  • 前者(RESTful)主要用於各組件之間的通信(如nova與glance的通信),或者說用於組件對外提供調用接口
  • 而後者(RPC)則用於同一組件中各個不同模塊之間的通信(如nova組件中nova-compute與nova-scheduler的通信)。

關於OpenStack中基於RESTful API的通信方式主要是應用了WSGI,這個知識點,我在前一篇文章中,有深入地講解過,你可以點擊查看。

對於不熟悉 OpenStack 的人,也別擔心聽不懂,這樣吧,我給你提兩個問題:

  1. RPC 和 REST 區別是什麼?
  2. 為什麼要採用RPC呢?

第一個問題:RPC 和 REST 區別是什麼?

你一定會覺得這個問題很奇怪,是的,包括我,但是你在網絡上一搜,會發現類似對比的文章比比皆是,我在想可能很多初學者由於基礎不牢固,才會將不相干的二者拿出來對比吧。既然是這樣,那為了讓你更加了解陌生的RPC,就從你熟悉得不能再熟悉的 REST 入手吧。

01、所屬類別不同

REST,是Representational State Transfer 的簡寫,中文描述表述性狀態傳遞(是指某個瞬間狀態的資源數據的快照,包括資源數據的內容、表述格式(XML、JSON)等信息。)

REST 是一種軟件架構風格。 這種風格的典型應用,就是HTTP。其因為簡單、擴展性強的特點而廣受開發者的青睞。

而RPC 呢,是 Remote Procedure Call Protocol 的簡寫,中文描述是遠程過程調用,它可以實現客戶端像調用本地服務(方法)一樣調用服務器的服務(方法)。

RPC 是一種基於 TCP 的通信協議,按理說它和REST不是一個層面上的東西,不應該放在一起討論,但是誰讓REST這麼流行呢,它是目前最流行的一套互聯網應用程序的API設計標準,某種意義下,我們說 REST 可以其實就是指代 HTTP 協議。

02、使用方式不同

從使用上來看,HTTP 接口只關注服務提供方,對於客戶端怎麼調用並不關心。接口只要保證有客戶端調用時,返回對應的數據就行了。而RPC則要求客戶端接口保持和服務端的一致。

  • REST 是服務端把方法寫好,客戶端並不知道具體方法。客戶端只想獲取資源,所以發起HTTP請求,而服務端接收到請求后根據URI經過一系列的路由才定位到方法上面去
  • PRC是服務端提供好方法給客戶端調用,客戶端需要知道服務端的具體類,具體方法,然後像調用本地方法一樣直接調用它。

03、面向對象不同

從設計上來看,RPC,所謂的遠程過程調用 ,是面向方法的 ,REST:所謂的 Representational state transfer ,是面向資源的,除此之外,還有一種叫做 SOA,所謂的面向服務的架構,它是面向消息的,這個接觸不多,就不多說了。

04、序列化協議不同

接口調用通常包含兩個部分,序列化和通信協議。

通信協議,上面已經提及了,REST 是 基於 HTTP 協議,而 RPC 可以基於 TCP/UDP,也可以基於 HTTP 協議進行傳輸的。

常見的序列化協議,有:json、xml、hession、protobuf、thrift、text、bytes等,REST 通常使用的是 JSON或者XML,而 RPC 使用的是 JSON-RPC,或者 XML-RPC。

通過以上幾點,我們知道了 REST 和 RPC 之間有很明顯的差異。

第二個問題:為什麼要採用RPC呢?

那到底為何要使用 RPC,單純的依靠RESTful API不可以嗎?為什麼要搞這麼多複雜的協議,渣渣表示真的學不過來了。

關於這一點,以下幾點僅是我的個人猜想,僅供交流哈:

  1. RPC 和 REST 兩者的定位不同,REST 面向資源,更注重接口的規範,因為要保證通用性更強,所以對外最好通過 REST。而 RPC 面向方法,主要用於函數方法的調用,可以適合更複雜通信需求的場景。
  2. RESTful API客戶端與服務端之間採用的是同步機制,當發送HTTP請求時,客戶端需要等待服務端的響應。當然對於這一點是可以通過一些技術來實現異步的機制的。
  3. 採用RESTful API,客戶端與服務端之間雖然可以獨立開發,但還是存在耦合。比如,客戶端在發送請求的時,必須知道服務器的地址,且必須保證服務器正常工作。而 rpc + ralbbimq中間件可以實現低耦合的分佈式集群架構。

說了這麼多,我們該如何選擇這兩者呢?我總結了如下兩點,供你參考:

  • REST 接口更加規範,通用適配性要求高,建議對外的接口都統一成 REST(也有例外,比如我接觸過 zabbix,其 API 就是基於 JSON-RPC 2.0協議的)。而組件內部的各個模塊,可以選擇 RPC,一個是不用耗費太多精力去開發和維護多套的HTTP接口,一個RPC的調用性能更高(見下條)
  • 從性能角度看,由於HTTP本身提供了豐富的狀態功能與擴展功能,但也正由於HTTP提供的功能過多,導致在網絡傳輸時,需要攜帶的信息更多,從性能角度上講,較為低效。而RPC服務網絡傳輸上僅傳輸與業務內容相關的數據,傳輸數據更小,性能更高。

02. 實現遠程調用的三種方式

“遠程調用”意思就是:被調用方法的具體實現不在程序運行本地,而是在別的某個地方(分佈到各個服務器),調用者只想要函數運算的結果,卻不需要實現函數的具體細節。

01、基於 xml-rpc

Python實現 rpc,可以使用標準庫里的 SimpleXMLRPCServer,它是基於XML-RPC 協議的。

有了這個模塊,開啟一個 rpc server,就變得相當簡單了。執行以下代碼:

import SimpleXMLRPCServer

class calculate:
    def add(self, x, y):
        return x + y

    def multiply(self, x, y):
        return x * y

    def subtract(self, x, y):
        return abs(x-y)

    def divide(self, x, y):
        return x/y


obj = calculate()
server = SimpleXMLRPCServer.SimpleXMLRPCServer(("localhost", 8088))
# 將實例註冊給rpc server
server.register_instance(obj)

print "Listening on port 8088"
server.serve_forever()

有了 rpc server,接下來就是 rpc client,由於我們上面使用的是 XML-RPC,所以 rpc clinet 需要使用xmlrpclib 這個庫。

import xmlrpclib

server = xmlrpclib.ServerProxy("http://localhost:8088")

然後,我們通過 server_proxy 對象就可以遠程調用之前的rpc server的函數了。

>> server.add(2, 3)
5
>>> server.multiply(2, 3)
6
>>> server.subtract(2, 3)
1
>>> server.divide(2, 3)
0

SimpleXMLRPCServer是一個單線程的服務器。這意味着,如果幾個客戶端同時發出多個請求,其它的請求就必須等待第一個請求完成以後才能繼續。

若非要使用 SimpleXMLRPCServer 實現多線程併發,其實也不難。只要將代碼改成如下即可。

from SimpleXMLRPCServer import SimpleXMLRPCServer
from SocketServer import ThreadingMixIn
class ThreadXMLRPCServer(ThreadingMixIn, SimpleXMLRPCServer):pass

class MyObject:
    def hello(self):
        return "hello xmlprc"

obj = MyObject()
server = ThreadXMLRPCServer(("localhost", 8088), allow_none=True)
server.register_instance(obj)

print "Listening on port 8088"
server.serve_forever()

02、基於json-rpc

SimpleXMLRPCServer 是基於 xml-rpc 實現的遠程調用,上面我們也提到 除了 xml-rpc 之外,還有 json-rpc 協議。

那 python 如何實現基於 json-rpc 協議呢?

答案是很多,很多web框架其自身都自己實現了json-rpc,但我們要獨立這些框架之外,要尋求一種較為乾淨的解決方案,我查找到的選擇有兩種

第一種是 jsonrpclib

pip install jsonrpclib -i https://pypi.douban.com/simple

第二種是 python-jsonrpc

pip install python-jsonrpc -i https://pypi.douban.com/simple

先來看第一種 jsonrpclib

它與 Python 標準庫的 SimpleXMLRPCServer 很類似(因為它的類名就叫做 SimpleJSONRPCServer ,不明真相的人真以為它們是親兄弟)。或許可以說,jsonrpclib 就是仿照 SimpleXMLRPCServer 標準庫來進行編寫的。

它的導入與 SimpleXMLRPCServer 略有不同,因為SimpleJSONRPCServer分佈在jsonrpclib庫中。

服務端

from jsonrpclib.SimpleJSONRPCServer import SimpleJSONRPCServer

server = SimpleJSONRPCServer(('localhost', 8080))
server.register_function(lambda x,y: x+y, 'add')
server.serve_forever()

客戶端

import jsonrpclib

server = jsonrpclib.Server("http://localhost:8080")

再來看第二種python-jsonrpc,寫起來貌似有些複雜。

服務端

import pyjsonrpc


class RequestHandler(pyjsonrpc.HttpRequestHandler):

    @pyjsonrpc.rpcmethod
    def add(self, a, b):
        """Test method"""
        return a + b

http_server = pyjsonrpc.ThreadingHttpServer(
    server_address=('localhost', 8080),
    RequestHandlerClass=RequestHandler
)
print "Starting HTTP server ..."
print "URL: http://localhost:8080"
http_server.serve_forever()

客戶端

import pyjsonrpc

http_client = pyjsonrpc.HttpClient(
    url="http://localhost:8080/jsonrpc"
)

還記得上面我提到過的 zabbix API,因為我有接觸過,所以也拎出來講講。zabbix API 也是基於 json-rpc 2.0協議實現的。

因為內容較多,這裏只帶大家打個,zabbix 是如何調用的:直接指明要調用 zabbix server 的哪個方法,要傳給這個方法的參數有哪些。

03、基於 zerorpc

以上介紹的兩種rpc遠程調用方式,如果你足夠細心,可以發現他們都是http+rpc 兩種協議結合實現的。

接下來,我們要介紹的這種(zerorpc),就不再使用走 http 了。

zerorpc 這個第三方庫,它是基於TCP協議、 ZeroMQ 和 MessagePack的,速度相對快,響應時間短,併發高。zerorpc 和 pyjsonrpc 一樣,需要額外安裝,雖然SimpleXMLRPCServer不需要額外安裝,但是SimpleXMLRPCServer性能相對差一些。

pip install zerorpc -i https://pypi.douban.com/simple

服務端代碼

import zerorpc

class caculate(object):
    def hello(self, name):
        return 'hello, {}'.format(name)

    def add(self, x, y):
        return x + y

    def multiply(self, x, y):
        return x * y

    def subtract(self, x, y):
        return abs(x-y)

    def divide(self, x, y):
        return x/y

s = zerorpc.Server(caculate())

s.bind("tcp://0.0.0.0:4242")
s.run()

客戶端

import zerorpc

c = zerorpc.Client()
c.connect("tcp://127.0.0.1:4242")

客戶端除了可以使用zerorpc框架實現代碼調用之外,它還支持使用“命令行”的方式調用。

客戶端可以使用命令行,那服務端是不是也可以呢?

是的,通過 Github 上的文檔幾個 demo 可以體驗到這個第三方庫做真的是優秀。

比如我們可以用下面這個命令,創建一個rpc server,後面這個 time Python 標準庫中的 time 模塊,zerorpc 會將 time 註冊綁定以供client調用。

zerorpc --server --bind tcp://127.0.0.1:1234 time

在客戶端,就可以用這條命令來遠程調用這個 time 函數。

zerorpc --client --connect tcp://127.0.0.1:1234 strftime %Y/%m/%d

03. 往rpc中引入消息中間件

經過了上面的學習,我們已經學會了如何使用多種方式實現rpc遠程調用。

通過對比,zerorpc 可以說是脫穎而出,一支獨秀。

但為何在 OpenStack 中,rpc client 不直接 rpc 調用 rpc server ,而是先把 rpc 調用請求發給 RabbitMQ ,再由訂閱者(rpc server)來取消息,最終實現遠程調用呢?

為此,我也做了一番思考:

OpenStack 組件繁多,在一個較大的集群內部每個組件內部通過rpc通信頻繁,如果都採用rpc直連調用的方式,連接數會非常地多,開銷大,若有些 server 是單線程的模式,超時會非常的嚴重。

OpenStack 是複雜的分佈式集群架構,會有多個 rpc server 同時工作,假設有 server01,server02,server03 三個server,當 rpc client 要發出rpc請求時,發給哪個好呢?這是問題一。

你可能會說輪循或者隨機,這樣對大家都公平。這樣的話還會引出另一個問題,倘若請求剛好發到server01,而server01剛好不湊巧,可能由於機器或者其他因為導致服務沒在工作,那這個rpc消息可就直接失敗了呀。要知道做為一個集群,高可用是基本要求,如果出現剛剛那樣的情況其實是很尷尬的。這是問題二。

集群有可能根據實際需要擴充節點數量,如果使用直接調用,耦合度太高,不利於部署和生產。這是問題三。

引入消息中間件,可以很好的解決這些問題。

解決問題一:消息只有一份,接收者由AMQP的負載算法決定,默認為在所有Receiver中均勻發送(round robin)。

解決問題二:有了消息中間件做緩衝站,client 可以任性隨意的發,server 都掛掉了?沒有關係,等 server 正常工作后,自己來消息中間件取就行了。

解決問題三:無論有多少節點,它們只要認識消息中間件這一个中介就足夠了。

04. 消息隊列你應該知道什麼?

由於後面,我將實例講解 OpenStack 中如何將 rpc 和 mq broker 結合使用。

而在此之前,你必須對消息隊列的一些基本知識有個概念。

首先,RPC只是定義了一個通信接口,其底層的實現可以各不相同,可以是 socket,也可以是今天要講的 AMQP。

AMQP(Advanced Message Queuing Protocol)是一種基於隊列的可靠消息服務協議,作為一種通信協議,AMQP同樣存在多個實現,如Apache Qpid,RabbitMQ等。

以下是 AMQP 中的幾個必知的概念:

  • Publisher:消息發布者

  • Receiver:消息接收者,在RabbitMQ中叫訂閱者:Subscriber。

  • Queue:用來保存消息的存儲空間,消息沒有被receiver前,保存在隊列中。

  • Exchange:用來接收Publisher發出的消息,根據Routing key 轉發消息到對應的Message Queue中,至於轉到哪個隊列里,這個路由算法又由exchange type決定的。

    exchange type:主要四種描述exchange的類型。

    direct:消息路由到滿足此條件的隊列中(queue,可以有多個): routing key = binding key

    topic:消息路由到滿足此條件的隊列中(queue,可以有多個):routing key 匹配 binding pattern. binding pattern是類似正則表達式的字符串,可以滿足複雜的路由條件。

    fanout:消息路由到多有綁定到該exchange的隊列中。

  • binding :binding是用來描述exchange和queue之間的關係的概念,一個exchang可以綁定多個隊列,這些關係由binding建立。前面說的binding key /binding pattern也是在binding中給出。

在網上找了個圖,可以很清晰地描述幾個名詞的關係。

關於AMQP,有幾下幾點值得注意:

  1. 每個receiver/subscriber 在接收消息前都需要創建binding。
  2. 一個隊列可以有多個receiver,隊列里的一個消息只能發給一個receiver。
  3. 一個消息可以被發送到一個隊列中,也可以被發送到多個多列中。多隊列情況下,一個消息可以被多個receiver收到並處理。Openstack RPC中這兩種情況都會用到。

05. OpenStack中如何使用RPC?

前面鋪墊了那麼久,終於到了講真實應用的場景。在生產中RPC是如何應用的呢?

其他模型我不太清楚,在 OpenStack 中的應用模型是這樣的

至於為什麼要如此設計,前面我已經給出了自己的觀點。

接下來,就是源碼解讀 OpenStack ,看看其是如何通過rpc進行遠程調用的。如若你對此沒有興趣(我知道很多人對此都沒有興趣,所以不浪費大家時間),可以直接跳過這一節,進入下一節。

目前Openstack中有兩種RPC實現,一種是在oslo messaging,一種是在openstack.common.rpc。

openstack.common.rpc是舊的實現,oslo messaging是對openstack.common.rpc的重構。openstack.common.rpc在每個項目中都存在一份拷貝,oslo messaging即將這些公共代碼抽取出來,形成一個新的項目。oslo messaging也對RPC API 進行了重新設計,對多種 transport 做了進一步封裝,底層也是用到了kombu這個AMQP庫。(注:Kombu 是Python中的messaging庫。Kombu旨在通過為AMQ協議提供慣用的高級接口,使Python中的消息傳遞盡可能簡單,併為常見的消息傳遞問題提供經過驗證和測試的解決方案。)

關於oslo_messaging庫,主要提供了兩種獨立的API:

  1. oslo.messaging.rpc(實現了客戶端-服務器遠程過程調用)
  2. oslo.messaging.notify(實現了事件的通知機制)

因為 notify 實現是太簡單了,所以這裏我就不多說了,如果有人想要看這方面內容,可以收藏我的博客(http://python-online.cn) ,我會更新補充 notify 的內容。

OpenStack RPC 模塊提供了 rpc.call,rpc.cast, rpc.fanout_cast 三種 RPC 調用方法,發送和接收 RPC 請求。

  • rpc.call 發送 RPC 同步請求並返回請求處理結果。
  • rpc.cast 發送 RPC 異步請求,與 rpc.call 不同之處在於,不需要請求處理結果的返回。
  • rpc.fanout_cast 用於發送 RPC 廣播信息無返回結果

rpc.call 和 rpc.rpc.cast 從實現代碼上看,他們的區別很小,就是call調用時候會帶有wait_for_reply=True參數,而cast不帶。

要了解 rpc 的調用機制呢,首先要知道 oslo_messaging 的幾個概念

  • transport:RPC功能的底層實現方法,這裡是rabbitmq的消息隊列的訪問路徑

    transport 就是定義你如何訪連接消息中間件,比如你使用的是 Rabbitmq,那在 nova.conf中應該有一行transport_url的配置,可以很清楚地看出指定了 rabbitmq 為消息中間件,並配置了連接rabbitmq的user,passwd,主機,端口。

    transport_url=rabbit://user:passwd@host:5672

    def get_transport(conf, url=None, allowed_remote_exmods=None):
        return _get_transport(conf, url, allowed_remote_exmods,
                              transport_cls=RPCTransport)
  • target:指定RPC topic交換機的匹配信息和綁定主機。

    target用來表述 RPC 服務器監聽topic,server名稱和server監聽的exchange,是否廣播fanout。

    class Target(object):
            def __init__(self, exchange=None, topic=None, namespace=None,
                     version=None, server=None, fanout=None,
                     legacy_namespaces=None):
            self.exchange = exchange
            self.topic = topic
            self.namespace = namespace
            self.version = version
            self.server = server
            self.fanout = fanout
            self.accepted_namespaces = [namespace] + (legacy_namespaces or [])

    rpc server 要獲取消息,需要定義target,就像一個門牌號一樣。

    rpc client 要發送消息,也需要有target,說明消息要發到哪去。

  • endpoints:是可供別人遠程調用的對象

    RPC服務器暴露出endpoint,每個 endpoint 包涵一系列的可被遠程客戶端通過 transport 調用的方法。直觀理解,可以參考nova-conductor創建rpc server的代碼,這邊的endpoints就是 nova/manager.py:ConductorManager()

  • dispatcher:分發器,這是 rpc server 才有的概念 只有通過它 server 端才知道接收到的rpc請求,要交給誰處理,怎麼處理?

    在client端,是這樣指定要調用哪個方法的。

    而在server端,是如何知道要執行這個方法的呢?這就是dispatcher 要乾的事,它從 endpoint 里找到這個方法,然後執行,最後返回。

  • Serializer:在 python 對象和message(notification) 之間數據做序列化或是反序列化的基類。

    主要方法有四個:

    1. deserialize_context(ctxt) :對字典變成 request contenxt.
    2. deserialize_entity(ctxt, entity) :對entity做反序列化,其中ctxt是已經deserialize過的,entity是要處理的。
    3. serialize_context(ctxt) :將Request context變成字典類型
    4. serialize_entity(ctxt, entity) :對entity做序列化,其中ctxt是已經deserialize過的,entity是要處理的。
  • executor:服務的運行方式,單線程或者多線程

    每個notification listener都和一個executor綁定,來控制收到的notification如何分配。默認情況下,使用的是blocking executor(具體特性參加executor一節)

    oslo_messaging.get_notification_listener(transport, targets, endpoints, executor=’blocking’, serializer=None, allow_requeue=False, pool=None)

rpc server 和rpc client 的四個重要方法

  1. reset():Reset service.
  2. start():該方法調用后,server開始poll,從transport中接收message,然後轉發給dispatcher.該message處理過程一直進行,直到stop方法被調用。executor決定server的IO處理策略。可能會是用一個新進程、新協程來做poll操作,或是直接簡單的在一個循環中註冊一個回調。同樣,executor也決定分配message的方式,是在一個新線程中dispatch或是….. *
  3. stop():當調用stop之後,新的message不會被處理。但是,server可能還在處理一些之前沒有處理完的message,並且底層driver資源也還一直沒有釋放。
  4. wait():在stop調用之後,可能還有message正在被處理,使用wait方法來阻塞當前進程,直到所有的message都處理完成。之後,底層的driver資源會釋放。

06. 模仿OpenStack寫rpc調用

模仿是一種很高效的學習方法,我這裏根據 OpenStack 的調用方式,抽取出核心內容,寫成一個簡單的 demo,有對 OpenStack 感興趣的可以了解一下,大部分人也可以直接跳過這章節

以下代碼不能直接運行,你還需要配置 rabbitmq 的連接方式,你可以寫在配置文件中,通過 get_transport 從cfg.CONF 中讀取,也可以直接將其寫成url的格式做成參數,傳給 get_transport 。

簡單的 rpc client

#coding=utf-8
import oslo_messaging
from oslo_config import cfg

# 創建 rpc client
transport = oslo_messaging.get_transport(cfg.CONF, url="")
target = oslo_messaging.Target(topic='test', version='2.0')
client = oslo_messaging.RPCClient(transport, target)

# rpc同步調用
client.call(ctxt, 'test', arg=arg)

簡單的 rpc server

#coding=utf-8
from oslo_config import cfg
import oslo_messaging
import time

# 定義endpoint類
class ServerControlEndpoint(object):
    target = oslo_messaging.Target(namespace='control',
                                   version='2.0')

    def __init__(self, server):
        self.server = server

    def stop(self, ctx):
        if self.server:
            self.server.stop()

            
class TestEndpoint(object):

    def test(self, ctx, arg):
        return arg

    
# 創建rpc server
transport = oslo_messaging.get_transport(cfg.CONF, url="")
target = oslo_messaging.Target(topic='test', server='server1')
endpoints = [
    ServerControlEndpoint(None),
    TestEndpoint(),
]
server = oslo_messaging.get_rpc_server(transport, target,endpoints,executor='blocking')
try:
    server.start()
    while True:
        time.sleep(1)
except KeyboardInterrupt:
    print("Stopping server")

server.stop()
server.wait()

以上,就是本期推送的全部內容,周末两天沒有出門,都在寫這篇文章。真的快掏空了我自己,不過寫完后真的很暢快。

【精選推薦文章】

智慧手機時代的來臨,RWD網頁設計已成為網頁設計推薦首選

想知道網站建置、網站改版該如何進行嗎?將由專業工程師為您規劃客製化網頁設計及後台網頁設計

帶您來看台北網站建置台北網頁設計,各種案例分享

廣告預算用在刀口上,網站設計公司幫您達到更多曝光效益