首页 文章 精选 留言 我的

精选列表

搜索[源码学习],共10000篇文章
优秀的个人博客,低调大师

hbase region split源码分析

hbase region split : split执行调用流程: 1.HbaseAdmin发起split:### 2.RSRpcServices实现类执行split(Implements the regionserver RPC services.)### 3.CompactSplitThread类与SplitRequest类用来执行region切割:### 4.splitRequest执行doSplitting操作### 4.1初始化两个子region### 4.2执行切割#### 4.2.1:(创建子region。stepsBeforePONR函数)##### 4.2.2:子region执行:stepsAfterPONR函数执行(包含openDaughters函数):##### 4.2.3:HRegionServer添加子region到meta表,加入RegionServer##### 4.3等待region切分完成,修改meta表信息,报告master#### 1.HbaseAdmin发起split: public void split(final TableName tableName) throws IOException { split(tableName, null); } public void split(final ServerName sn, final HRegionInfo hri, byte[] splitPoint) throws IOException { if (hri.getStartKey() != null && splitPoint != null && Bytes.compareTo(hri.getStartKey(), splitPoint) == 0) { throw new IOException("should not give a splitkey which equals to startkey!"); } // TODO: There is no timeout on this controller. Set one! HBaseRpcController controller = rpcControllerFactory.newController(); controller.setPriority(hri.getTable()); // TODO: this does not do retries, it should. Set priority and timeout in controller AdminService.BlockingInterface admin = this.connection.getAdmin(sn); //hbase中split调用请求通过protobuf(AdminProtos)实现 ProtobufUtil.split(controller, admin, hri, splitPoint); 实现代码: try { admin.splitRegion(controller, request); } 2.RSRpcServices实现类执行split(Implements the regionserver RPC services.) 代码: try { checkOpen(); requestCount.increment(); Region region = getRegion(request.getRegion()); region.startRegionOperation(Operation.SPLIT_REGION); if (region.getRegionInfo().getReplicaId() != HRegionInfo.DEFAULT_REPLICA_ID) { throw new IOException("Can't split replicas directly. " + "Replicas are auto-split when their primary is split."); } LOG.info("Splitting " + region.getRegionInfo().getRegionNameAsString()); region.flush(true); byte[] splitPoint = null; if (request.hasSplitPoint()) { //确定切割点splitpoint splitPoint = request.getSplitPoint().toByteArray(); } ((HRegion)region).forceSplit(splitPoint); //请求region切割 regionServer.compactSplitThread.requestSplit(region, ((HRegion)region).checkSplit(), RpcServer.getRequestUser()); return SplitRegionResponse.newBuilder().build(); } 3.CompactSplitThread类与SplitRequest类用来执行region切割: 代码: try { this.splits.execute(new SplitRequest(r, midKey, this.server, user)); if (LOG.isDebugEnabled()) { LOG.debug("Split requested for " + r + ". " + this); } } try { //acquire a shared read lock on the table, so that table schema modifications //do not happen concurrently //获取table的读锁。。 tableLock = server.getTableLockManager().readLock(parent.getTableDesc().getTableName() , "SPLIT_REGION:" + parent.getRegionInfo().getRegionNameAsString()); try { tableLock.acquire(); } 4.splitRequest执行doSplitting操作 4.1### //初始化两个子region信息 this.hri_a = new HRegionInfo(hri.getTable(), startKey, this.splitrow, false, rid); this.hri_b = new HRegionInfo(hri.getTable(), this.splitrow, endKey, false, rid); 4.2### //执行切割 st.execute(this.server, this.server, user); 4.2.1:(创建子region。stepsBeforePONR函数)#### //创建两个子region PairOfSameType<Region> regions = createDaughters(server, services, user); 4.2.2:stepsAfterPONR函数执行(openDaughters):#### // 两个子region DaughterOpener线程 start DaughterOpener aOpener = new DaughterOpener(server, (HRegion)a); DaughterOpener bOpener = new DaughterOpener(server, (HRegion)b); //向hdfs上写入.regionInfo文件以便meta挂掉以便恢复 writeRegionInfoOnFilesystem(content, true); //初始化所有的hstore initializeStores(reporter, status); //LoadStoreFiles函数执行: if (files == null || files.size() == 0) { return new ArrayList<StoreFile>(); } // initialize the thread pool for opening store files in parallel.. ThreadPoolExecutor storeFileOpenerThreadPool = this.region.getStoreFileOpenAndCloseThreadPool("StoreFileOpenerThread-" + this.getColumnFamilyName()); CompletionService<StoreFile> completionService = new ExecutorCompletionService<StoreFile>(storeFileOpenerThreadPool); int totalValidStoreFile = 0; for (final StoreFileInfo storeFileInfo: files) { //HDFS上对应的路径和文件 completionService.submit(new Callable<StoreFile>() { @Override public StoreFile call() throws IOException { //每个文件创建一个StoreFile对象,对每个storefile对象会读取文件上的内容创建一个 HalfStoreFileReader读对象来操作该region的父region上的相应的文件,及该 region上目前存储的是引用文件,其指向的是其父region上的相应的文件,对该 region的所有读或写都将关联到父region上 StoreFile storeFile = createStoreFileAndReader(storeFileInfo); return storeFile; } }); totalValidStoreFile++; } //将子Region添加到rs的online region列表上,并添加到meta表上 { if (useZKForAssignment) { // add 2nd daughter first (see HBASE-4335) services.postOpenDeployTasks(b); } else if (!services.reportRegionStateTransition(TransitionCode.SPLIT, parent.getRegionInfo(), hri_a, hri_b)) { throw new IOException("Failed to report split region to master: " + parent.getRegionInfo().getShortNameToLog()); } // Should add it to OnlineRegions services.addToOnlineRegions(b); if (useZKForAssignment) { services.postOpenDeployTasks(a); } services.addToOnlineRegions(a); } 参考文章:http://blog.csdn.net/Pun_C/article/details/47173453

优秀的个人博客,低调大师

Storm-源码分析- metric

首先定义一系列metric相关的interface, IMetric, IReducer, ICombiner (backtype.storm.metric.api) 在task中, 创建一系列builtin-metrics, (backtype.storm.daemon.builtin-metrics), 并注册到topology context里面 task会不断的利用如spout-acked-tuple!的functions去更新这些builtin-metrics task会定期将builtin-metrics里面的统计数据通过METRICS-STREAM发送给metric-bolt (backtype.storm.metric.MetricsConsumerBolt, 该bolt会创建实现backtype.storm.metric.api.IMetricsConsumer的对象, 用于计算出metrics) 然后如何使用这些metrics? 由于这是builtin metrics, 是不会被外界使用的 如果处理这些metrics, 取决于_metricsConsumer.handleDataPoints, 这里的_metricsConsumer是通过topology's configuration配置的 比如backtype.storm.metric.LoggingMetricsConsumer, 如果使用这个consumer就会将metrics写入log中 1. backtype.storm.metric.api IMetric package backtype.storm.metric.api; public interface IMetric { public Object getValueAndReset(); ;;取得当前值并恢复初始状态 } CountMetric, 计数, reset时清零 AssignableMetric, 赋值, 不用reset MultiCountMetric, 使用hashmap记录多个count, reset时分别对每个count对象调用getValueAndReset public class CountMetric implements IMetric { long _value = 0; public CountMetric() { } public void incr() { _value++; } public void incrBy(long incrementBy) { _value += incrementBy; } public Object getValueAndReset() { long ret = _value; _value = 0; return ret; } } ICombiner public interface ICombiner<T> { public T identity(); public T combine(T a, T b); } CombinedMetric, 结合ICombiner和IMetric public class CombinedMetric implements IMetric { private final ICombiner _combiner; private Object _value; public CombinedMetric(ICombiner combiner) { _combiner = combiner; _value = _combiner.identity(); } public void update(Object value) { _value = _combiner.combine(_value, value); } public Object getValueAndReset() { Object ret = _value; _value = _combiner.identity(); return ret; } } IReducer publicinterfaceIReducer<T> { T init(); T reduce(T accumulator, Object input); Object extractResult(T accumulator); } 实现IReducer接口, 实现平均数Reducer, reduce里面累加和计数, extractResult里面acc/count求平均数 class MeanReducerState { public int count = 0; public double sum = 0.0; } public class MeanReducer implements IReducer<MeanReducerState> { public MeanReducerState init() { return new MeanReducerState(); } public MeanReducerState reduce(MeanReducerState acc, Object input) { acc.count++; if(input instanceof Double) { acc.sum += (Double)input; } else if(input instanceof Long) { acc.sum += ((Long)input).doubleValue(); } else if(input instanceof Integer) { acc.sum += ((Integer)input).doubleValue(); } else { throw new RuntimeException( "MeanReducer::reduce called with unsupported input type `" + input.getClass() + "`. Supported types are Double, Long, Integer."); } return acc; } public Object extractResult(MeanReducerState acc) { if(acc.count > 0) { return new Double(acc.sum / (double)acc.count); } else { return null; } } } ReducedMetric 结合IReducer和IMetric public class ReducedMetric implements IMetric { private final IReducer _reducer; private Object _accumulator; public ReducedMetric(IReducer reducer) { _reducer = reducer; _accumulator = _reducer.init(); } public void update(Object value) { _accumulator = _reducer.reduce(_accumulator, value); } public Object getValueAndReset() { Object ret = _reducer.extractResult(_accumulator); _accumulator = _reducer.init(); return ret; } } IMetricsConsumer 这个interface, 内嵌TaskInfo和DataPoint类 handleDataPoints, 添加逻辑以处理task对应的一系列DataPoint public interface IMetricsConsumer { public static class TaskInfo { public TaskInfo() {} public TaskInfo(String srcWorkerHost, int srcWorkerPort, String srcComponentId, int srcTaskId, long timestamp, int updateIntervalSecs) { this.srcWorkerHost = srcWorkerHost; this.srcWorkerPort = srcWorkerPort; this.srcComponentId = srcComponentId; this.srcTaskId = srcTaskId; this.timestamp = timestamp; this.updateIntervalSecs = updateIntervalSecs; } public String srcWorkerHost; public int srcWorkerPort; public String srcComponentId; public int srcTaskId; public long timestamp; public int updateIntervalSecs; } public static class DataPoint { public DataPoint() {} public DataPoint(String name, Object value) { this.name = name; this.value = value; } @Override public String toString() { return "[" + name + " = " + value + "]"; } public String name; public Object value; } void prepare(Map stormConf, Object registrationArgument, TopologyContext context, IErrorReporter errorReporter); void handleDataPoints(TaskInfo taskInfo, Collection<DataPoint> dataPoints); void cleanup(); } 2. backtype.storm.daemon.builtin-metrics 定义Spout和Bolt所需要的一些metric, 主要两个record, BuiltinSpoutMetrics和BuiltinBoltMetrics, [metric-name, metric-object]的hashmap (defrecord BuiltinSpoutMetrics [^MultiCountMetric ack-count ^MultiReducedMetric complete-latency ^MultiCountMetric fail-count ^MultiCountMetric emit-count ^MultiCountMetric transfer-count]) (defrecord BuiltinBoltMetrics [^MultiCountMetric ack-count ^MultiReducedMetric process-latency ^MultiCountMetric fail-count ^MultiCountMetric execute-count ^MultiReducedMetric execute-latency ^MultiCountMetric emit-count ^MultiCountMetric transfer-count]) (defn make-data [executor-type] (condp = executor-type :spout (BuiltinSpoutMetrics. (MultiCountMetric.) (MultiReducedMetric. (MeanReducer.)) (MultiCountMetric.) (MultiCountMetric.) (MultiCountMetric.)) :bolt (BuiltinBoltMetrics. (MultiCountMetric.) (MultiReducedMetric. (MeanReducer.)) (MultiCountMetric.) (MultiCountMetric.) (MultiReducedMetric. (MeanReducer.)) (MultiCountMetric.) (MultiCountMetric.)))) (defn register-all [builtin-metrics storm-conf topology-context] (doseq [[kw imetric] builtin-metrics] (.registerMetric topology-context (str "__" (name kw)) imetric (int (get storm-conf Config/TOPOLOGY_BUILTIN_METRICS_BUCKET_SIZE_SECS))))) 在mk-task-data的时候, 调用make-data来创建相应的metrics, :builtin-metrics (builtin-metrics/make-data (:type executor-data)) 并在executor的mk-threads中, 会将这些builtin-metrics注册到topologycontext中去, (builtin-metrics/register-all (:builtin-metrics task-data) storm-conf (:user-context task-data)) 上面完成的builtin-metrics的创建和注册, 接着定义了一系列用于更新metrics的functions, 以spout-acked-tuple!为例, 需要更新MultiCountMetric ack-count和MultiReducedMetric complete-latency .scope从MultiCountMetric取出某个CountMetric, 然后incrBy来将stats的rate增加到count上 (defn spout-acked-tuple! [^BuiltinSpoutMetrics m stats stream latency-ms] (-> m .ack-count (.scope stream) (.incrBy (stats-rate stats))) (-> m .complete-latency (.scope stream) (.update latency-ms))) 3. backtype.storm.metric MetricsConsumerBolt 创建实现IMetricsConsumer的对象, 并在execute里面调用handleDataPoints package backtype.storm.metric; public class MetricsConsumerBolt implements IBolt { IMetricsConsumer _metricsConsumer; String _consumerClassName; OutputCollector _collector; Object _registrationArgument; public MetricsConsumerBolt(String consumerClassName, Object registrationArgument) { _consumerClassName = consumerClassName; _registrationArgument = registrationArgument; } @Override public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) { try { _metricsConsumer = (IMetricsConsumer)Class.forName(_consumerClassName).newInstance(); } catch (Exception e) { throw new RuntimeException("Could not instantiate a class listed in config under section " + Config.TOPOLOGY_METRICS_CONSUMER_REGISTER + " with fully qualified name " + _consumerClassName, e); } _metricsConsumer.prepare(stormConf, _registrationArgument, context, (IErrorReporter)collector); _collector = collector; } @Override public void execute(Tuple input) { _metricsConsumer.handleDataPoints((IMetricsConsumer.TaskInfo)input.getValue(0), (Collection)input.getValue(1)); _collector.ack(input); } @Override public void cleanup() { _metricsConsumer.cleanup(); } } SystemBolt SystemBolt, 根据comments里面说的, 每个worker都有一个, taskid=-1 定义些system相关的metric, 并注册到topologycontext里面 需要使用Java调用clojure, 所以需要import下面的package import clojure.lang.AFn; import clojure.lang.IFn; //funtion import clojure.lang.RT; //run-time 并且用到些用于监控memory和JVM的java package java.lang.management.MemoryUsage, 表示内存使用量快照的MemoryUsage对象 java.lang.management.GarbageCollectorMXBean, 用于Java虚拟机的垃圾回收的管理接口, 比如发生的回收的总次数, 和累计回收时间 java.lang.management.RuntimeMXBean, 用于Java 虚拟机的运行时系统的管理接口 这个bolt的特点是, 只有prepare实现了逻辑, 并且通过_prepareWasCalled保证prepare只被执行一次 prepare中的逻辑, 主要就是定义各种metric, 并且通过registerMetric注册到TopologyContext中 metic包含, JVM的运行时间, 开始时间, memory情况, 和每个GarbageCollector的情况 注册的这些system metrics也会一起被发送到MetricsConsumerBolt进行处理 这应该用spout实现, 为啥用bolt实现? // There is one task inside one executor for each worker of the topology. // TaskID is always -1, therefore you can only send-unanchored tuples to co-located SystemBolt. // This bolt was conceived to export worker stats via metrics api. public class SystemBolt implements IBolt { private static Logger LOG = LoggerFactory.getLogger(SystemBolt.class); private static boolean _prepareWasCalled = false; private static class MemoryUsageMetric implements IMetric { IFn _getUsage; public MemoryUsageMetric(IFn getUsage) { _getUsage = getUsage; } @Override public Object getValueAndReset() { MemoryUsage memUsage = (MemoryUsage)_getUsage.invoke(); HashMap m = new HashMap(); m.put("maxBytes", memUsage.getMax()); m.put("committedBytes", memUsage.getCommitted()); m.put("initBytes", memUsage.getInit()); m.put("usedBytes", memUsage.getUsed()); m.put("virtualFreeBytes", memUsage.getMax() - memUsage.getUsed()); m.put("unusedBytes", memUsage.getCommitted() - memUsage.getUsed()); return m; } } // canonically the metrics data exported is time bucketed when doing counts. // convert the absolute values here into time buckets. private static class GarbageCollectorMetric implements IMetric { GarbageCollectorMXBean _gcBean; Long _collectionCount; Long _collectionTime; public GarbageCollectorMetric(GarbageCollectorMXBean gcBean) { _gcBean = gcBean; } @Override public Object getValueAndReset() { Long collectionCountP = _gcBean.getCollectionCount(); Long collectionTimeP = _gcBean.getCollectionTime(); Map ret = null; if(_collectionCount!=null && _collectionTime!=null) { ret = new HashMap(); ret.put("count", collectionCountP - _collectionCount); ret.put("timeMs", collectionTimeP - _collectionTime); } _collectionCount = collectionCountP; _collectionTime = collectionTimeP; return ret; } } @Override public void prepare(final Map stormConf, TopologyContext context, OutputCollector collector) { if(_prepareWasCalled && stormConf.get(Config.STORM_CLUSTER_MODE) != "local") { throw new RuntimeException("A single worker should have 1 SystemBolt instance."); } _prepareWasCalled = true; int bucketSize = RT.intCast(stormConf.get(Config.TOPOLOGY_BUILTIN_METRICS_BUCKET_SIZE_SECS)); final RuntimeMXBean jvmRT = ManagementFactory.getRuntimeMXBean(); context.registerMetric("uptimeSecs", new IMetric() { @Override public Object getValueAndReset() { return jvmRT.getUptime()/1000.0; } }, bucketSize); context.registerMetric("startTimeSecs", new IMetric() { @Override public Object getValueAndReset() { return jvmRT.getStartTime()/1000.0; } }, bucketSize); context.registerMetric("newWorkerEvent", new IMetric() { boolean doEvent = true; @Override public Object getValueAndReset() { if (doEvent) { doEvent = false; return 1; } else return 0; } }, bucketSize); final MemoryMXBean jvmMemRT = ManagementFactory.getMemoryMXBean(); context.registerMetric("memory/heap", new MemoryUsageMetric(new AFn() { public Object invoke() { return jvmMemRT.getHeapMemoryUsage(); } }), bucketSize); context.registerMetric("memory/nonHeap", new MemoryUsageMetric(new AFn() { public Object invoke() { return jvmMemRT.getNonHeapMemoryUsage(); } }), bucketSize); for(GarbageCollectorMXBean b : ManagementFactory.getGarbageCollectorMXBeans()) { context.registerMetric("GC/" + b.getName().replaceAll("\\W", ""), new GarbageCollectorMetric(b), bucketSize); } } @Override public void execute(Tuple input) { throw new RuntimeException("Non-system tuples should never be sent to __system bolt."); } @Override public void cleanup() { } } 4. system-topology! 这里会动态的往topology里面, 加入metric-component (MetricsConsumerBolt) 和system-component (SystemBolt), 以及相应的steam信息 system-topology!会往topology加上些东西 1. acker, 后面再说 2. metric-bolt, input是所有component的tasks发来的METRICS-STREAM, 没有output 3. system-bolt, 没有input, output是两个TICK-STREAM 4. 给所有component, 增加额外的输出metrics-stream, system-stream (defn system-topology! [storm-conf ^StormTopology topology] (validate-basic! topology) (let [ret (.deepCopy topology)] (add-acker! storm-conf ret) (add-metric-components! storm-conf ret) (add-system-components! storm-conf ret) (add-metric-streams! ret) (add-system-streams! ret) (validate-structure! ret) ret )) 4.1 增加component 看下thrift中的定义, 往topology里面增加一个blot component, 其实就是往hashmap中增加一组[string, Bolt] 关键就是看看如何使用thrift/mk-bolt-spec*来创建blot spec struct StormTopology { 1: required map<string, SpoutSpec> spouts; 2: required map<string, Bolt> bolts; 3: required map<string, StateSpoutSpec> state_spouts; } struct Bolt { 1: required ComponentObject bolt_object; 2: required ComponentCommon common; } struct ComponentCommon { 1: required map<GlobalStreamId, Grouping> inputs; 2: required map<string, StreamInfo> streams; //key is stream id, outputs 3: optional i32 parallelism_hint; //how many threads across the cluster should be dedicated to this component 4: optional string json_conf; } struct StreamInfo { 1: required list<string> output_fields; 2: required bool direct; } (defn add-metric-components! [storm-conf ^StormTopology topology] (doseq [[comp-id bolt-spec] (metrics-consumer-bolt-specs storm-conf topology)] ;;从metrics-consumer-bolt-specs中可以看出该bolt会以METRICS-STREAM-ID为输入, 且没有输出 (.put_to_bolts topology comp-id bolt-spec))) (defn add-system-components! [conf ^StormTopology topology] (let [system-bolt-spec (thrift/mk-bolt-spec* {} ;;input为空, 没有输入 (SystemBolt.) ;;object {SYSTEM-TICK-STREAM-ID (thrift/output-fields ["rate_secs"]) METRICS-TICK-STREAM-ID (thrift/output-fields ["interval"])} ;;output, 定义两个output streams, 但代码中并没有emit :p 0 :conf {TOPOLOGY-TASKS 0})] (.put_to_bolts topology SYSTEM-COMPONENT-ID system-bolt-spec))) metric-components 首先, topology里面所有的component(包含system component), 都需要往metics-bolt发送统计数据, 所以component-ids-that-emit-metrics就是all-components-ids+SYSTEM-COMPONENT-ID 那么对于任意一个comp, 都会对metics-bolt产生如下输入, {[comp-id METRICS-STREAM-ID] :shuffle} (采用:suffle grouping方式) 然后, 用thrift/mk-bolt-spec*来定义创建bolt的fn, mk-bolt-spec 最后, 调用mk-bolt-spec来创建metics-bolt的spec, 参考上面的定义 关键就是, 创建MetricsConsumerBolt对象, 需要从storm-conf里面读出, MetricsConsumer的实现类和参数 这个bolt负责, 将从各个task接收到的数据, 调用handleDataPoints生成metircs, 参考前面的定义 (defn metrics-consumer-bolt-specs [storm-conf topology] (let [component-ids-that-emit-metrics (cons SYSTEM-COMPONENT-ID (keys (all-components topology))) inputs (->> (for [comp-id component-ids-that-emit-metrics] {[comp-id METRICS-STREAM-ID] :shuffle}) (into {})) mk-bolt-spec (fn [class arg p] (thrift/mk-bolt-spec* inputs ;;inputs集合 (backtype.storm.metric.MetricsConsumerBolt. class arg) ;;object {} ;;output为空 :p p :conf {TOPOLOGY-TASKS p}))] (map (fn [component-id register] [component-id (mk-bolt-spec (get register "class") (get register "argument") (or (get register "parallelism.hint") 1))]) (metrics-consumer-register-ids storm-conf) (get storm-conf TOPOLOGY-METRICS-CONSUMER-REGISTER)))) 4.2 增加stream 给每个component增加两个output stream METRICS-STREAM-ID, 发送给metric-blot, 数据结构为output-fields ["task-info" "data-points"] SYSTEM-STREAM-ID, ,数据结构为output-fields ["event"] (defn add-metric-streams! [^StormTopology topology] (doseq [[_ component] (all-components topology) :let [common (.get_common component)]] (.put_to_streams common METRICS-STREAM-ID (thrift/output-fields ["task-info" "data-points"])))) (defn add-system-streams! [^StormTopology topology] (doseq [[_ component] (all-components topology) :let [common (.get_common component)]] (.put_to_streams common SYSTEM-STREAM-ID (thrift/output-fields ["event"])))) 本文章摘自博客园,原文发布日期:2013-07-30

优秀的个人博客,低调大师

twitter storm源码走读(一)

nimbus启动场景分析 本文详细介绍了twitter storm中的nimbus节点的启动场景,分析nimbus是如何一步步实现定义于storm.thrift中的service,以及如何利用curator来和zookeeper server建立通讯。 对于storm client来说,nimbus是storm cluster与外部的唯一接口,是总的接口人,在这个接口上使用thrift定义的各种service。但是nimbus光接单并不干活,具体的脏活累活,这哥们都是分配到各个slots上的,让nimbus来具体管理各个slots也就是worker,似乎还是太累了,中层干部supervisor同学适时参与了。 nimbus并不知道到底有哪些supervisor会加入到自己的团队中,它啥时规定了每个supervisor最多能带几个worker。对于supervisor的加入与退出,是通过zookeeper server来告知的。好了,在下面的分析中,每个接口上的初始化工作具体有哪些将一一呈现。 tuple消息发送场景分析 worker进程内消息接收与处理全景图 先上幅图简要勾勒出worker进程接收到tuple消息之后的处理全过程 IConnection的建立与使用 话说在mk-threads :bolt函数的实现中有这么一段代码,其主要功能是实现tuple的emit功能 bolt-emit (fn [stream anchors values task] (let [out-tasks (if task (tasks-fn task stream values) (tasks-fn stream values))] (fast-list-iter [t out-tasks] (let [anchors-to-ids (HashMap.)] (fast-list-iter [^TupleImpl a anchors] (let [root-ids (-> a .getMessageId .getAnchorsToIds .keySet)] (when (pos? (count root-ids)) (let [edge-id (MessageId/generateId rand)] (.updateAckVal a edge-id) (fast-list-iter [root-id root-ids] (put-xor! anchors-to-ids root-id edge-id)) )))) (transfer-fn t (TupleImpl. worker-context values task-id stream (MessageId/makeId anchors-to-ids))))) (or out-tasks []))) 加亮为蓝色的部分实现的功能是另外发送tuple,那么transfer-fn函数的定义在哪呢?见mk-threads的let部分,能见到下述一行代码 :transfer-fn (mk-executor-transfer-fn batch-transfer->worker) 在继续往下看每个函数实现之前,先确定一下这节代码阅读的目的。storm在线程之间使用disruptor进行通讯,在进程之间进行消息通讯使用的是zeromq或netty, 所以需要从transfer-fn追踪到使用zeromq或netty api的位置。 再看mk-executor-transfer-fn函数实现 (defn mk-executor-transfer-fn [batch-transfer->worker] (fn this ([task tuple block? ^List overflow-buffer] (if (and overflow-buffer (not (.isEmpty overflow-buffer))) (.add overflow-buffer [task tuple]) (try-cause (disruptor/publish batch-transfer->worker [task tuple] block?) (catch InsufficientCapacityException e (if overflow-buffer (.add overflow-buffer [task tuple]) (throw e)) )))) ([task tuple overflow-buffer] (this task tuple (nil? overflow-buffer) overflow-buffer)) ([task tuple] (this task tuple nil) ))) disruptor/publish表示将消息从本线程发送出去,至于谁是该消息的接收者,请继续往下看。 worker进程中,有一个receiver-thread是用来专门接收来自外部进程的消息,那么与之相对的是有一个transfer-thread用来将本进程的消息发送给外部进程。所以刚才的disruptor/publish发送出来的消息应该被transfer-thread接收到。 在transfer-thread中,能找到这行下述一行代码 transfer-thread (disruptor/consume-loop* (:transfer-queue worker) transfer-tuples) 对于接收到来自本进程中其它线程发送过来的消息利用transfer-tuples进行处理,transfer-tuples使用mk-transfer-tuples-handler来创建,所以需要看看mk-transfer-tuples-handler能否与zeromq或netty联系上呢? (defn mk-transfer-tuples-handler [worker] (let [^DisruptorQueue transfer-queue (:transfer-queue worker) drainer (ArrayList.) node+port->socket (:cached-node+port->socket worker) task->node+port (:cached-task->node+port worker) endpoint-socket-lock (:endpoint-socket-lock worker) ] (disruptor/clojure-handler (fn [packets _ batch-end?] (.addAll drainer packets) (when batch-end? (read-locked endpoint-socket-lock (let [node+port->socket @node+port->socket task->node+port @task->node+port] ;; consider doing some automatic batching here (would need to not be serialized at this point to remo ;; try using multipart messages ... first sort the tuples by the target node (without changing the lo 17 (fast-list-iter [[task ser-tuple] drainer] ;; TODO: consider write a batch of tuples here to every target worker ;; group by node+port, do multipart send (let [node-port (get task->node+port task)] (when node-port (.send ^IConnection (get node+port->socket node-port) task ser-tuple)) )))) (.clear drainer)))))) 上述代码中出现了与zeromq可能有联系的部分了即加亮为红色的一行。 那凭什么说加亮的IConnection一行与zeromq有关系的,这话得慢慢说起,需要从配置文件开始。 在storm.yaml中有这么一行配置项,即 storm.messaging.transport: "backtype.storm.messaging.zmq" 这个配置项与worker中的mqcontext相对应,所以在worker中以mqcontext为线索,就能够一步步找到IConnection的实现。connections在函数mk-refresh-connections中建立 refresh-connections (mk-refresh-connections worker) mk-refresh-connection函数中与mq-context相关联的一部分代码如下所示 (swap! (:cached-node+port->socket worker) #(HashMap. (merge (into {} %1) %2)) (into {} (dofor [endpoint-str new-connections :let [[node port] (string->endpoint endpoint-str)]] [endpoint-str (.connect ^IContext (:mq-context worker) storm-id ((:node->host assignment) node) port) ] ))) 注意加亮部分,利用mq-conext中connect函数来创建IConnection. 当打开zmq.clj时候,就能验证我们的猜测。 (^IConnection connect [this ^String storm-id ^String host ^int port] (require 'backtype.storm.messaging.zmq) (-> context (mq/socket mq/push) (mq/set-hwm hwm) (mq/set-linger linger-ms) (mq/connect (get-connect-zmq-url local? host port)) mk-connection)) 代码走到这里,IConnection什么时候建立起来的谜底就揭开了,消息是如何从bolt或spout线程传递到transfer-thread,再由zeromq将tuple发送给下跳的路径打通了。 tuple的分发策略 grouping 从一个bolt中产生的tuple可以有多个bolt接收,到底发送给哪一个bolt呢?这牵扯到分发策略问题,其实在twitter storm中有两个层面的分发策略问题,一个是对于task level的,在讲topology submit的时候已经涉及到。另一个就是现在要讨论的针对tuple level的分发。 再次将视线拉回到bolt-emit中,这次将目光集中在变量t的前前后后。 (let [out-tasks (if task (tasks-fn task stream values) (tasks-fn stream values))] (fast-list-iter [t out-tasks] (let [anchors-to-ids (HashMap.)] (fast-list-iter [^TupleImpl a anchors] (let [root-ids (-> a .getMessageId .getAnchorsToIds .keySet)] (when (pos? (count root-ids)) (let [edge-id (MessageId/generateId rand)] (.updateAckVal a edge-id) (fast-list-iter [root-id root-ids] (put-xor! anchors-to-ids root-id edge-id)) )))) (transfer-fn t (TupleImpl. worker-context values task-id stream (MessageId/makeId anchors-to-ids))))) 上述代码显示t从out-tasks来,而out-tasks是tasks-fn的返回值 tasks-fn (:tasks-fn task-data) 一谈tasks-fn,原来从未涉及的文件task.clj这次被挂上了,task-data与由task/mk-task创建。将中间环节跳过,调用关系如下所列。 mk-task mk-task-data mk-tasks-fn tasks-fn中会使用到grouping,处理代码如下 fn ([^Integer out-task-id ^String stream ^List values] (when debug? (log-message "Emitting direct: " out-task-id "; " component-id " " stream " " values)) (let [target-component (.getComponentId worker-context out-task-id) component->grouping (get stream->component->grouper stream) grouping (get component->grouping target-component) out-task-id (if grouping out-task-id)] (when (and (not-nil? grouping) (not= :direct grouping)) (throw (IllegalArgumentException. "Cannot emitDirect to a task expecting a regular grouping"))) (apply-hooks user-context .emit (EmitInfo. values stream task-id [out-task-id])) (when (emit-sampler) (builtin-metrics/emitted-tuple! (:builtin-metrics task-data) executor-stats stream) (stats/emitted-tuple! executor-stats stream) (if out-task-id (stats/transferred-tuples! executor-stats stream 1) (builtin-metrics/transferred-tuple! (:builtin-metrics task-data) executor-stats stream 1))) (if out-task-id [out-task-id]) )) 而每个topology中的grouping策略又是如何被executor知道的呢,这从另一端executor-data说起。 在mk-executor-data中有下面一行代码 :stream->component->grouper (outbound-components worker-context component-id) outbound-components的定义如下 (defn outbound-components "Returns map of stream id to component id to grouper" [^WorkerTopologyContext worker-context component-id] (->> (.getTargets worker-context component-id) clojurify-structure (map (fn [[stream-id component->grouping]] [stream-id (outbound-groupings worker-context component-id stream-id (.getComponentOutputFields worker-context component-id stream-id) component->grouping)])) (into {}) (HashMap.)))

优秀的个人博客,低调大师

twitter storm源码走读(四)

Trident Topology执行过程分析 TridentTopology是storm提供的高层使用接口,常见的一些SQL中的操作在tridenttopology提供的api中都有类似的影射。关于TridentTopology的使用及运行原理,当前进行详细分析的文章不多。 从TridentTopology到vanilla topology(普通的topology)由三个层次组成: 面向最终用户的概念stream, operation 利用planner将tridenttopology转换成vanilla topology 执行vanilla topology 本文尝试TridentTopology是如何先一步步转换成普通的storm Topology(即vanila topology), 转换后的topology的执行中有哪些区别? 概述 从TridentTopology到基本的Topology有三层,下图给出一个全局的视图。 创建TridentTopology 下面的代码摘自StormStarter中的TridentWordCount.java TridentTopology topology = new TridentTopology(); topology.newStream("spout1", spout).parallelismHint(16).each(new Fields("sentence"), new Split(), new Fields("word")).groupBy(new Fields("word")).persistentAggregate(new MemoryMapState.Factory(), new Count(), new Fields("count")).parallelismHint(16); return topology.build(); 上述代码的newStream一行,分两大部分,一是使用newStream来创建一个stream对象,然后针对该Stream进行各种操作,each/shuffle/persistentAggregate等就是各种operation. 用户在使用TridentTopology的时候,只需要熟悉Stream和TridentTopology中的API函数即可。 转换TridentTopology为Vanilla Topology 上一节创建了Stream,但是如何将其与原有的Spout及Bolt联系起来呢?问题的关键就在TridentTopology::build函数和TridentTopologyBuilder::buildTopology TridentTopology::build newStream及其后的函数调用创建了一个含有三大类节点的List,利用该List创建了一个有向非循环图(DAG)。这三类节点分别是operation, partition, spout,在build函数将节点分类分别加入到boltNodes或spoutNodes,注意此处的spout或bolt不能等同于普通的spout和bolt. TridentTopologyBuilder::buildTopology 利用在build函数中创建的boltNodes,spoutNodes及生成的graph来创建vanilla topology所需要的bolt及spout. 在buildTopology中会看到类似的代码片段。 builder.setBolt(spoutCoordinator(id), new TridentSpoutCoordinator(c.commitStateId, (ITridentSpout) c.spout)) .globalGrouping(masterCoordinator(c.batchGroupId), MasterBatchCoordinator.BATCH_STREAM_ID) .globalGrouping(masterCoordinator(c.batchGroupId), MasterBatchCoordinator.SUCCESS_STREAM_ID); builder.setSpout(masterCoordinator(batch), new MasterBatchCoordinator(commitIds, batchesToSpouts.get(batch))); for(String b: c.committerBatches) { specs.get(b).commitStream = new GlobalStreamId(masterCoordinator(b), MasterBatchCoordinator.COMMIT_STREAM_ID); } BoltDeclarer d = builder.setBolt(id, new TridentBoltExecutor(c.bolt, batchIdsForBolts, specs), c.parallelism); 最终生成的普通Topology,与普通Topology中的Spout相对应的是MasterBatchCoordinator,而在创建TridentTopology使用的spout则成了Bolt,使用于Stream上的各种Operation也存在于多个普通Bolt中。 TridentTopology的执行 TridentTopology被转换为普通的Topology(vanilla Topology)之后提交到nimbus,它的具体执行过程有什么不同呢? 主要有几点: MasterBatchCoordinator通过Batch_stream_id来发送通知给TridentSpoutExecutor TridentSpoutExecutor收到通知发送成批的tuple给下一跳的Bolt 下一跳的Bolt收到tuple之后,使用TridentBoltExecutor来进行处理 TridentBoltExecutor调用SubtopologyBolt::execute InitialReceiver::execute被调用 TridentProcessor::execute被调用 MasterBatchCoordinator收到ack之后,会发送success消息给Spout MasterBatchCoordinator在commit的时候,会发送commit消息给Spout,让Spout将缓存的消息删除 trident topology可靠性分析 本文详细分析TridentTopology的可靠性实现, TridentTopology通过transactional spout与transactional state相结合,能够做到tuple“只被处理一次,不多也不少”。也就是做到事务性处理exactly-once,要么成功,要么失败。 而一般的storm topology是无法保证eactly-once的处理的,它们要么是at-least-once(至少被处理一次,有可能被处理多次);要么是at-most-once(最多被处理一次,这样就存在遗漏的可能). TridentTopology在设计中借鉴和保留了目前已经过期的transactional topology的设计思想。 Storm Topology的ack机制 在进行TridentTopology的可靠性分析之前,我们先回顾一下在storm topology中的ack机制。ack bolt是在提交到storm cluster中,由系统自动产生的,一般来说一个topology只有一个ack bolt(当然可以通过配置参数指定多个)。 当bolt处理并下发完tuple给下一跳的bolt时,会发送一个ack给ack bolt。ack bolt通过简单的异或原理(即同一个数与自己异或结果为零)来判定从spout发出的某一个bolt是否已经被完全处理完毕。如果结果为真,ack bolt发送消息给spout,spout中的ack函数被调用并执行。如果超时,则发送fail消息给spout,spout中的fail函数被调用并执行,spout中的ack和fail的处理逻辑由用户自行填写。 如在github上的kerstel spout就能做到只有当某一个tuple被成功处理之后,它才会从缓存中移除,否则继续放入到处理队列再次进行处理。 TridentTopology的可靠性机制 在“走读之6”一文中分析了一个tridenttopology是如何转换成storm topology的,我想用上面这幅图再次阐述一下转变后的结果。 一个tridenttopoloy会至少引入一个MasterBatchCoordinator,这个MBC就类似于storm topology中的spout newStream时使用的入参spout会裂变成两个bolt,一是TridentSpoutCoordinator,另一个是TridentSpoutExecutor 针对stream的各种操作则被分散到各个Bolt中,它们的执行上下文是TridentBoltExecutor 可以看出使用TridentTopology Api进行操作时,所有的东西其实都运行在bolt context中,而真正的spout是在调用TridentTopologyBuilder.buildTopology()的时候被添加的。 MasterBatchCoordinator使用batch_stream发送一个类似于seeder tuple的东西给tridentspoutcoordinator,tridentspoutcoordinator将该信号继续下发给TridentSpoutExecutor, TridentSpout是如何一步步被调用到的呢。 TridentBoltExecutor::execute TridentSpoutExecutor::execute BatchSpoutExecutor::execute ITridentSpout::emitBatch emitBatch是产生真正需要被处理的tuple的,这些tuple会被各个Operation所在的bolt所接收。它们的调用顺序是 TridentBoltExecutor::execute SubtopologyBolt::execute InitialReceiver::receive TridentProcessor::execute 处理结束的判断依据 在TridentSpout中是如何判断所有的tuple都已经被处理的呢。 在每跳中认为自己处理完毕的时候,它都会告诉下一跳,即下游,我给你发送了多少tuple,如果下游将上游发送过来的确认消息与自身确实已经处理的消息比对一致的话,则认为处理都完成,于是发送ack. 问题的关键变成每一个bolt是如何判断自己已经处理完毕的呢,请看步骤3 总有一个bolt是没有上游的,即TridentSpoutExecutor,它只会收到启动指令,但不接收真正的业务数据,于是它会告诉下一跳,我发了多少tuple给你。 STREAM 在MasterBatchCoordinator中定义了三种不同的stream,这三种stream分别是 BATCH_STREAM COMMIT_STREAM SUCCESS_STREAM 这些stream分别在什么时候被使用呢,下图给出一个大概的时序 简要说明: masterbatchcoordinator通过batch_stream发送seeder tuple给tridentspoutcoordinator tridentspoutcoordinator给tridentspoutexecutor继续传递该指令 TridentSpoutExecutor在收到启动指令后,调用ITridentSpout接口的实现类进行emitBatch TridentSpoutExecutor在发送完一批batch后,finishBatch被调用,通过emitDirect会给下一跳通过coord_stream发送trackedinfo,即我已经发送了多少消息给你 TridentSpoutExecutor紧接着还会给ack bolt发送ack消息,ack bolt将其传达到MasterBatchCoordinator MasterBatchCoordinator在收到第一个ack后,将状态置为processed 当MasterBatchCoordinator再次收到ack后,会将状态转为committing,同时通过commit_stream发送tuple给TridentSpoutExecutor 收到commit_stream上传来的tuple后,TridentSpoutExecutor会调用ITridentSpout中的emmitter, emmitter::commit()被执行,TridentSpoutExecutor会再次ack收到tuple MasterBatchCoordinator在收到这个tuple之后,会认为针对某一个seeder tuple的处理已经完全实现,于是通过SUCCESS_STREM告知TridentSpoutCoordinator,所有的活都已经都完成了,收工。 收到Success_stream上传来的信号后,ITridentSpout中的内嵌子类Emmit和Coordinator中相应的success方法会被调用执行。 注意: 为了描述方便,将TridentTopology进行了简化,认为其在转换成真正的storm topology时,只有一个TridentProcessor所在的bolt。真实的情况可能比这复杂,但消息的传递路径还是差不多的。 注意在TridentTopology中ack会被多次反复调用,这不同于普通的storm topology 状态机 在MasterBatchCoordinator中,针对每一个seeder tuple,其状态机如下图所示。注意这些状态是会被保存到zookeeper server中的,使用的api定义在TransactionalState中。 总结 通过上面的分析可以看出,TridentTopology实现了一个比较好的框架,但真正要做到exactly-once的处理,还需要用户自己去实现ITridentSpout中的两个重要内嵌类,Emmitter和Coordinator。 具体如何实现该接口,可以查看storm-core/src/jvm/storm/trident/testing目录下的FixedBatchSpout.java和FeederCommitterBatchSpout.java

优秀的个人博客,低调大师

twitter storm源码走读(五)

TridentTopology创建过程详解 从用户层面来看TridentTopology,有两个重要的概念一是Stream,另一个是作用于Stream上的各种Operation。在实现层面来看,无论是stream,还是后续的operation都会转变成为各个Node,这些Node之间的关系通过重要的数据结构图来维护。具体到TridentTopology,实现图的各种操作的组件是jgrapht。 说到图,两个基本的概念会闪现出来,一是结点,二是描述结点之间关系的边。要想很好的理解TridentTopology就需要紧盯图中结点和边的变化。 TridentTopology在转换成为普通的StormTopology时,需要将原始的图分成各个group,每个group将运行于一个独立的bolt中。TridentTopology又是如何知道哪些node应该在同一个group,哪些应该处在另一个group中的呢;如何来确定每个group的并发度(parallismHint)的呢。这些问题的解决都与jgrapht分不开。 关于jgrapht的更多信息,请参考其官方网站http://jgrapht.org 概要 在TridentTopology中向图中添加结点的api有三种: addNode addSourcedNode addSourcedStateNode 其中addNode在创建stream是使用,addSourcedStateNode在partitionPersist时使用到,其它的operation使用到的是addSourcedNode. addNode与其它两个方法的一个重要区别还在于,addNode是不需要添加边(Edge),而其它两个API需要往图中添加edge,以确定该node的源是哪个。 TridentTopology 1 2 3 4 public TridentTopology() { _graph = new DefaultDirectedGraph( new ErrorEdgeFactory()); _gen = new UniqueIdGen(); } 在TridentTopology的构造函数中,创建了DAG(有向无环图)。利用这个_graph来作为容器以存储后续过程中创建的各个node及它们之间的关系。 newStream newStream会为DAG(有向无环图)中创建源结点,其调用关系如下所示。 newStream addNode registerNode 1 protected void registerNode(Node n) { 2 _graph.addVertex(n); 3 if(n.stateInfo!=null) { 4 String id = n.stateInfo.id; 5 if(!_colocate.containsKey(id)) { 6 _colocate.put(id, new ArrayList()); 7 } 8 _colocate.get(id).add(n); 9 } 10 } each 作用于stream上的Operation有很多,以each为例来看新的operation是如何转换成为node添加到_graph中的。 //Stream.javapublic Stream each(Fields inputFields, Function function, Fields functionFields) { projectionValidation(inputFields); return _topology.addSourcedNode(this, new ProcessorNode(_topology.getUniqueStreamId(), _name, TridentUtils.fieldsConcat(getOutputFields(), functionFields), functionFields, new EachProcessor(inputFields, function))); } 调用关系描述如下 Stream::each TridentTopology::addSourcedNode TridentTopology::registerSourcedNode registerSourcedNode的实现如下 protected void registerSourcedNode(List<Stream> sources, Node newNode) { registerNode(newNode); int streamIndex = 0; for(Stream s: sources) { _graph.addEdge(s._node, newNode, new IndexedEdge(s._node, newNode, streamIndex)); streamIndex++; } } 注意此处添加edge是,是有索引的,这样可以区别处理的先后顺序。 在Stream中含有成员变量_node,表示stream最近停泊的node,有了该变量添加edge才成为了可能。 partitionPersist public TridentState partitionPersist(StateSpec stateSpec, Fields inputFields, StateUpdater updater, Fields functionFields) { projectionValidation(inputFields); String id = _topology.getUniqueStateId(); ProcessorNode n = new ProcessorNode(_topology.getUniqueStreamId(), _name, functionFields, functionFields, new PartitionPersistProcessor(id, inputFields, updater)); n.committer = true; n.stateInfo = new NodeStateInfo(id, stateSpec); return _topology.addSourcedStateNode(this, n); } 调用关系 Stream::partitionPersist TridentTopology::addSourcedStateNode TridentTopology::registerSourcedNode 与addNode及addSourcedNode不同的是,addSourcedStateNode返回的是TridentState而非Stream。 既然谈到了TridentState就不得不谈到其另一面Stream::stateQuery, public Stream stateQuery(TridentState state, Fields inputFields, QueryFunction function, Fields functionFields) { projectionValidation(inputFields); String stateId = state._node.stateInfo.id; Node n = new ProcessorNode(_topology.getUniqueStreamId(), _name, TridentUtils.fieldsConcat(getOutputFields(), functionFields), functionFields, new StateQueryProcessor(stateId, inputFields, function)); _topology._colocate.get(stateId).add(n); return _topology.addSourcedNode(this, n); } 从此处可以看出stateQueryNode最起码有两个inputStream,一是从TridentState而来表示状态已经改变,另一个是处于drpcStream这个方面的上一跳结点。 build TridentTopology::build是将TridentTopology转变为StormTopology的过程,这一过程中最重要的一环就是将_graph中含有的node进行分组。 grouping 算法逻辑概述 将boltNodes中的每个boltNode作为一个group加入全部加入initialGroups 以graph和initialGroups作为入参创建GraphGrouper 分组的过程其实就是进行合并的过程,详见GraphGrouper::mergeFully() 如果从当前group1的输出目的地都是属于group2,则将group1,group2合并 如果当前group1的所有输入源都是来自于group2,则将group1,group2合并 将需要合并的group1,group2作为入参创建新的group,同时将group1,group2从已有的集合出移除 public void mergeFully() { boolean somethingHappened = true; while(somethingHappened) { somethingHappened = false; for(Group g: currGroups) { Collection<Group> outgoingGroups = outgoingGroups(g); if(outgoingGroups.size()==1) { Group out = outgoingGroups.iterator().next(); if(out!=null) { merge(g, out); somethingHappened = true; break; } } Collection<Group> incomingGroups = incomingGroups(g); if(incomingGroups.size()==1) { Group in = incomingGroups.iterator().next(); if(in!=null) { merge(g, in); somethingHappened = true; break; } } } } } GraphGrouper::merge() private void merge(Group g1, Group g2) { Group newGroup = new Group(g1, g2); currGroups.remove(g1); currGroups.remove(g2); currGroups.add(newGroup); for(Node n: newGroup.nodes) { groupIndex.put(n, newGroup); } } 在group之间添加partitionNode // add identity partitions between groups for(IndexedEdge<Node> e: new HashSet<IndexedEdge>(graph.edgeSet())) { if(!(e.source instanceof PartitionNode) && !(e.target instanceof PartitionNode)) { Group g1 = grouper.nodeGroup(e.source); Group g2 = grouper.nodeGroup(e.target); // g1 being null means the source is a spout node if(g1==null && !(e.source instanceof SpoutNode)) throw new RuntimeException("Planner exception: Null source group must indicate a spout node at this phase of planning"); if(g1==null || !g1.equals(g2)) { graph.removeEdge(e); PartitionNode pNode = makeIdentityPartition(e.source); graph.addVertex(pNode); graph.addEdge(e.source, pNode, new IndexedEdge(e.source, pNode, 0)); graph.addEdge(pNode, e.target, new IndexedEdge(pNode, e.target, e.index)); } } } _graph中所有的node在变换过后,变成两组元素,一是spoutNodes,另一个是合并后的mergedGroup. spoutNodes中的每个元素作为spout添加到TridentTopologyBuilder的_spouts数组中,mergedGroup中的每个group添加到TridentTopologyBuilder的_bolt数组中。在TridentTopologyBuilder::build()中最主要的事情是为每个_spouts和_bolts数组中的成员添加grouping关系。 小结 到目前为止,通过两篇文章分析了TridentTopology的创建过程及其运行时在每个TridentBoltExecutor中的消息传递情况。接下来会分析TridentTopology提供的API实现及其作用场景。

优秀的个人博客,低调大师

Spark RDD类源码阅读

每天进步一点点~开搞~ abstract class RDD[T: ClassTag]( //@transient 注解表示将字段标记为瞬态的 @transient private var _sc: SparkContext, // Seq是序列,元素有插入的先后顺序,可以有重复的元素。 @transient private var deps: Seq[Dependency[_]] ) extends Serializable with Logging { if (classOf[RDD[_]].isAssignableFrom(elementClassTag.runtimeClass)) { user programs that } //这里应该是声明sparkContext对象后才能使用RDD的调用 private def sc: SparkContext = { if (_sc == null) { throw new SparkException( "RDD transformations and actions can only be invoked by the driver, not inside of other " + "transformations; for example, rdd1.map(x => rdd2.values.count() * x) is invalid because " + "the values transformation and count action cannot be performed inside of the rdd1.map " + "transformation. For more information, see SPARK-5063.") } _sc } //构建一个RDD应该是一对一的关系,比如子RDD对应唯一的父RDD def this(@transient oneParent: RDD[_]) = this(oneParent.context , List(new OneToOneDependency(oneParent))) private[spark] def conf: SparkConf = _conf //sparkconf的设置 def getConf: SparkConf = conf.clone() //获取相应的配置信息 def jars: Seq[String] = _jars def files: Seq[String] = _files def master: String = _conf.get("spark.master") def appName: String = _conf.get("spark.app.name") private[spark] def isEventLogEnabled: Boolean = _conf.getBoolean("spark.eventLog.enabled", false) private[spark] def eventLogDir: Option[URI] = _eventLogDir private[spark] def eventLogCodec: Option[String] = _eventLogCodec //临时文件夹的名称为spark+随机时间戳 val externalBlockStoreFolderName = "spark-" + randomUUID.toString() //判断是否为local模式 def isLocal: Boolean = (master == "local" || master.startsWith("local[")) //用于触发事件的监听 private[spark] val listenerBus = new LiveListenerBus // 该方法可用于测试用 private[spark] def createSparkEnv( conf: SparkConf, isLocal: Boolean, listenerBus: LiveListenerBus): SparkEnv = { SparkEnv.createDriverEnv(conf, isLocal, listenerBus, SparkContext.numDriverCores(master)) } //加载env配置文件 private[spark] def env: SparkEnv = _env private[spark] val addedFiles = HashMap[String, Long]() private[spark] val addedJars = HashMap[String, Long]() //监听所有调用persist的RDD private[spark] val persistentRdds = new TimeStampedWeakValueHashMap[Int, RDD[_]] //重用配置hadoop Configuration def hadoopConfiguration: Configuration = _hadoopConfiguration //用于设置executorMemory的内存数量 private[spark] def executorMemory: Int = _executorMemory // 将环境参数传递给exeuctor private[spark] val executorEnvs = HashMap[String, String]() // 设置正在使用SparkContext的用户 val sparkUser = Utils.getCurrentUserName() //设置提交的appliaction的唯一标识。就是当提交给yarn或local模式时,申请资源的applaction名称 def applicationId: String = _applicationId def applicationAttemptId: Option[String] = _applicationAttemptId def metricsSystem: MetricsSystem = if (_env != null) _env.metricsSystem else null private[spark] def eventLogger: Option[EventLoggingListener] = _eventLogger private[spark] def executorAllocationManager: Option[ExecutorAllocationManager] = _executorAllocationManager private[spark] def cleaner: Option[ContextCleaner] = _cleaner private[spark] var checkpointDir: Option[String] = None // 用户可以使用本地变量来传递消息 protected[spark] val localProperties = new InheritableThreadLocal[Properties] { override protected def childValue(parent: Properties): Properties = { //clone一下,防止父变量改变从而影响子变量 semantics (SPARK-10563). if (conf.get("spark.localProperties.clone", "false").toBoolean) { SerializationUtils.clone(parent).asInstanceOf[Properties] } else { new Properties(parent) } } override protected def initialValue(): Properties = new Properties() } private def warnSparkMem(value: String): String = { logWarning("Using SPARK_MEM to set amount of memory to use per executor process is " + "deprecated, please use spark.executor.memory instead.") value } //设置log级别,包括ALL, DEBUG, ERROR, FATAL, INFO, OFF, TRACE, WARN def setLogLevel(logLevel: String) { val validLevels = Seq("ALL", "DEBUG", "ERROR", "FATAL", "INFO", "OFF", "TRACE", "WARN") if (!validLevels.contains(logLevel)) { throw new IllegalArgumentException( s"Supplied level $logLevel did not match one of: ${validLevels.mkString(",")}") } Utils.setLogLevel(org.apache.log4j.Level.toLevel(logLevel)) } //不同模式的配置参数 if (!_conf.contains("spark.master")) { throw new SparkException("A master URL must be set in your configuration") } if (!_conf.contains("spark.app.name")) { throw new SparkException("An application name must be set in your configuration") } // System property spark.yarn.app.id must be set if user code ran by AM on a YARN cluster // yarn-standalone is deprecated, but still supported if ((master == "yarn-cluster" || master == "yarn-standalone") && !_conf.contains("spark.yarn.app.id")) { throw new SparkException("Detected yarn-cluster mode, but isn't running on a cluster. " + "Deployment to YARN is not supported directly by SparkContext. Please use spark-submit.") } _conf.setIfMissing("spark.driver.host", Utils.localHostName()) _conf.setIfMissing("spark.driver.port", "0") _conf.set("spark.executor.id", SparkContext.DRIVER_IDENTIFIER) _jars = _conf.getOption("spark.jars").map(_.split(",")).map(_.filter(_.size != 0)).toSeq.flatten _files = _conf.getOption("spark.files").map(_.split(",")).map(_.filter(_.size != 0)) .toSeq.flatten _eventLogDir = if (isEventLogEnabled) { val unresolvedDir = conf.get("spark.eventLog.dir", EventLoggingListener.DEFAULT_LOG_DIR) .stripSuffix("/") Some(Utils.resolveURI(unresolvedDir)) } else { None } _eventLogCodec = { val compress = _conf.getBoolean("spark.eventLog.compress", false) if (compress && isEventLogEnabled) { Some(CompressionCodec.getCodecName(_conf)).map(CompressionCodec.getShortName) } else { None } } //jobProgressListener应该在创建sparkEnv之前,因为当创建sparkEnv时,一些信息将会被发送到jobProgressListener,否则就会丢失啦。 _jobProgressListener = new JobProgressListener(_conf) listenerBus.addListener(jobProgressListener) _env = createSparkEnv(_conf, isLocal, listenerBus) SparkEnv.set(_env) _metadataCleaner = new MetadataCleaner(MetadataCleanerType.SPARK_CONTEXT, this.cleanup, _conf) _statusTracker = new SparkStatusTracker(this) _progressBar = if (_conf.getBoolean("spark.ui.showConsoleProgress", true) && !log.isInfoEnabled) { Some(new ConsoleProgressBar(this)) } else { None } _ui = if (conf.getBoolean("spark.ui.enabled", true)) { Some(SparkUI.createLiveUI(this, _conf, listenerBus, _jobProgressListener, _env.securityManager, appName, startTime = startTime)) } else { None } if (jars != null) { jars.foreach(addJar) } if (files != null) { files.foreach(addFile) } //获取启动app设置的参数变量,如果没有则获取配置文件中的 _executorMemory = _conf.getOption("spark.executor.memory") .orElse(Option(System.getenv("SPARK_EXECUTOR_MEMORY"))) .orElse(Option(System.getenv("SPARK_MEM")) .map(warnSparkMem)) .map(Utils.memoryStringToMb) .getOrElse(1024) //500这里在创建HeartbeatReceiver 之前先创建createTaskScheduler,因为每个Executor在构造函数中检索HeartbeatReceiver _heartbeatReceiver = env.rpcEnv.setupEndpoint( HeartbeatReceiver.ENDPOINT_NAME, new HeartbeatReceiver(this))

优秀的个人博客,低调大师

MapReduce源码分析之JobSplitWriter

JobSplitWriter被作业客户端用于写分片相关文件,包括分片数据文件job.split和分片元数据信息文件job.splitmetainfo。它有两个静态成员变量,如下: // 分片版本,当前默认为1 private static final int splitVersion = JobSplit.META_SPLIT_VERSION; // 分片文件头部,为UTF-8格式的字符串"SPL"的字节数组"SPL" private static final byte[] SPLIT_FILE_HEADER; 并且,提供了一个静态方法,完成SPLIT_FILE_HEADER的初始化,代码如下: // 静态方法,加载SPLIT_FILE_HEADER为UTF-8格式的字符串"SPL"的字节数组byte[] static { try { SPLIT_FILE_HEADER = "SPL".getBytes("UTF-8"); } catch (UnsupportedEncodingException u) { throw new RuntimeException(u); } } JobSplitWriter实现其功能的为createSplitFiles()方法,它有三种实现,我们先看其中的public static <T extends InputSplit> void createSplitFiles(Path jobSubmitDir,Configuration conf, FileSystem fs, T[] splits),代码如下: // 创建分片文件 public static <T extends InputSplit> void createSplitFiles(Path jobSubmitDir, Configuration conf, FileSystem fs, T[] splits) throws IOException, InterruptedException { // 调用createFile()方法,创建分片文件,并获取文件系统数据输出流FSDataOutputStream实例out, // 对应路径为jobSubmitDir/job.split,jobSubmitDir为参数yarn.app.mapreduce.am.staging-dir指定的路径/作业所属用户user/.staging/作业ID FSDataOutputStream out = createFile(fs, JobSubmissionFiles.getJobSplitFile(jobSubmitDir), conf); // 调用writeNewSplits()方法,将分片数据写入分片文件,并得到分片元数据信息SplitMetaInfo数组info SplitMetaInfo[] info = writeNewSplits(conf, splits, out); // 关闭输出流 out.close(); // 调用writeJobSplitMetaInfo()方法,将分片元数据信息写入分片元数据文件 writeJobSplitMetaInfo(fs,JobSubmissionFiles.getJobSplitMetaFile(jobSubmitDir), new FsPermission(JobSubmissionFiles.JOB_FILE_PERMISSION), splitVersion, info); } createSplitFiles()方法的逻辑很清晰,大体如下: 1、调用createFile()方法,创建分片文件,并获取文件系统数据输出流FSDataOutputStream实例out,对应路径为jobSubmitDir/job.split,jobSubmitDir为参数yarn.app.mapreduce.am.staging-dir指定的路径/作业所属用户user/.staging/作业ID; 2、调用writeNewSplits()方法,将分片数据写入分片文件,并得到分片元数据信息SplitMetaInfo数组info; 3、关闭输出流out; 4、调用writeJobSplitMetaInfo()方法,将分片元数据信息写入分片元数据文件。 我们先来看下createFile()方法,代码如下: private static FSDataOutputStream createFile(FileSystem fs, Path splitFile, Configuration job) throws IOException { // 调用HDFS文件系统FileSystem的create()方法,获取文件系统数据输出流FSDataOutputStream实例out, // 对应权限为JobSubmissionFiles.JOB_FILE_PERMISSION,即0644,rw-r--r-- FSDataOutputStream out = FileSystem.create(fs, splitFile, new FsPermission(JobSubmissionFiles.JOB_FILE_PERMISSION)); // 获取副本数replication,取参数mapreduce.client.submit.file.replication,参数未配置默认为10 int replication = job.getInt(Job.SUBMIT_REPLICATION, 10); // 通过文件系统FileSystem实例fs的setReplication()方法,设置splitFile的副本数位10 fs.setReplication(splitFile, (short)replication); // 调用writeSplitHeader()方法写入分片头信息 writeSplitHeader(out); // 返回文件系统数据输出流out return out; } 首先,调用HDFS文件系统FileSystem的create()方法,获取文件系统数据输出流FSDataOutputStream实例out,对应权限为JobSubmissionFiles.JOB_FILE_PERMISSION,即0644,rw-r--r--; 其次,获取副本数replication,取参数mapreduce.client.submit.file.replication,参数未配置默认为10; 接着,通过文件系统FileSystem实例fs的setReplication()方法,设置splitFile的副本数位10; 然后,调用writeSplitHeader()方法写入分片头信息; 最后,返回文件系统数据输出流out。 writeSplitHeader()方法专门用于将分片头部信息写入分片文件,代码如下: private static void writeSplitHeader(FSDataOutputStream out) throws IOException { // 文件系统数据输出流out写入byte[],内容为UTF-8格式的"SPL" out.write(SPLIT_FILE_HEADER); // 文件系统数据输出流out写入int,分片版本号,目前为1 out.writeInt(splitVersion); } 很简单,首先文件系统数据输出流out写入byte[],内容为UTF-8格式的"SPL",然后文件系统数据输出流out写入int,分片版本号,目前为1。 接下来,我们再看下writeNewSplits()方法,它将分片数据写入分片文件,并得到分片元数据信息SplitMetaInfo数组info,代码如下: @SuppressWarnings("unchecked") private static <T extends InputSplit> SplitMetaInfo[] writeNewSplits(Configuration conf, T[] array, FSDataOutputStream out) throws IOException, InterruptedException { // 根据array的大小,构造同等大小的分片元数据信息SplitMetaInfo数组info, // array其实是传入的分片数组 SplitMetaInfo[] info = new SplitMetaInfo[array.length]; if (array.length != 0) {// 如果array中有数据 // 创建序列化工厂SerializationFactory实例factory SerializationFactory factory = new SerializationFactory(conf); int i = 0; // 获取最大的数据块位置maxBlockLocations,取参数mapreduce.job.max.split.locations,参数未配置默认为10 int maxBlockLocations = conf.getInt(MRConfig.MAX_BLOCK_LOCATIONS_KEY, MRConfig.MAX_BLOCK_LOCATIONS_DEFAULT); // 通过输出流out的getPos()方法获取输出流out的当前位置offset long offset = out.getPos(); // 遍历数组array中每个元素split for(T split: array) { // 通过输出流out的getPos()方法获取输出流out的当前位置prevCount long prevCount = out.getPos(); // 往输出流out中写入String,内容为split对应的类名 Text.writeString(out, split.getClass().getName()); // 获取序列化器Serializer实例serializer Serializer<T> serializer = factory.getSerializer((Class<T>) split.getClass()); // 打开serializer,接入输出流out serializer.open(out); // 将split序列化到输出流out serializer.serialize(split); // 通过输出流out的getPos()方法获取输出流out的当前位置currCount long currCount = out.getPos(); // 通过split的getLocations()方法,获取位置信息locations String[] locations = split.getLocations(); if (locations.length > maxBlockLocations) { LOG.warn("Max block location exceeded for split: " + split + " splitsize: " + locations.length + " maxsize: " + maxBlockLocations); locations = Arrays.copyOf(locations, maxBlockLocations); } // 构造split对应的元数据信息,并加入info指定位置, // offset为当前split在split文件中的起始位置,数据长度为split.getLength(),位置信息为locations info[i++] = new JobSplit.SplitMetaInfo( locations, offset, split.getLength()); // offset增加当前split已写入数据大小 offset += currCount - prevCount; } } // 返回分片元数据信息SplitMetaInfo数组info return info; } writeNewSplits()方法的逻辑比较清晰,大体如下: 1、根据array的大小,构造同等大小的分片元数据信息SplitMetaInfo数组info,array其实是传入的分片数组; 2、如果array中有数据: 2.1、创建序列化工厂SerializationFactory实例factory; 2.2、获取最大的数据块位置maxBlockLocations,取参数mapreduce.job.max.split.locations,参数未配置默认为10; 2.3、通过输出流out的getPos()方法获取输出流out的当前位置offset; 2.4、遍历数组array中每个元素split: 2.4.1、通过输出流out的getPos()方法获取输出流out的当前位置prevCount; 2.4.2、往输出流out中写入String,内容为split对应的类名; 2.4.3、获取序列化器Serializer实例serializer; 2.4.4、打开serializer,接入输出流out; 2.4.5、将split序列化到输出流out; 2.4.6、通过输出流out的getPos()方法获取输出流out的当前位置currCount; 2.4.7、通过split的getLocations()方法,获取位置信息locations; 2.4.8、确保位置信息locations的长度不能超过maxBlockLocations,超过则截断; 2.4.9、构造split对应的元数据信息,并加入info指定位置,offset为当前split在split文件中的起始位置,数据长度为split.getLength(),位置信息为locations; 2.4.10、offset增加当前split已写入数据大小; 3、返回分片元数据信息SplitMetaInfo数组info。 其中,序列化split对象时,我们以FileSplit为例来分析,其write()方法如下: @Override public void write(DataOutput out) throws IOException { // 写入文件路径全名 Text.writeString(out, file.toString()); // 写入分片在文件中的起始位置 out.writeLong(start); // 写入分片在文件中的长度 out.writeLong(length); } 比较简单,分别写入文件路径全名、分片在文件中的起始位置、分片在文件中的长度三个信息。 综上所述,分片文件job.split文件的内容为: 1、文件头:"SPL"+int类型版本号1; 2、分片类信息:String类型split对应类名; 3、分片数据信息:String类型文件路径全名+Long类型分片在文件中的起始位置+Long类型分片在文件中的长度。 而在最后,构造分片元数据信息时,产生的是JobSplit的静态内部类SplitMetaInfo对象,包括分片位置信息locations、split在split文件中的起始位置offset、分片长度split.getLength()。 下面,我们再看下分片的元数据信息文件是如何产生的,让我们来研究下writeJobSplitMetaInfo()方法,代码如下: // 写入作业分片元数据信息 private static void writeJobSplitMetaInfo(FileSystem fs, Path filename, FsPermission p, int splitMetaInfoVersion, JobSplit.SplitMetaInfo[] allSplitMetaInfo) throws IOException { // write the splits meta-info to a file for the job tracker // 调用HDFS文件系统FileSystem的create()方法,生成分片元数据信息文件,并获取文件系统数据输出流FSDataOutputStream实例out, // 对应文件路径为jobSubmitDir/job.splitmetainfo,jobSubmitDir为参数yarn.app.mapreduce.am.staging-dir指定的路径/作业所属用户user/.staging/作业ID // 对应权限为JobSubmissionFiles.JOB_FILE_PERMISSION,即0644,rw-r--r-- FSDataOutputStream out = FileSystem.create(fs, filename, p); // 写入分片元数据头部信息UTF-8格式的字符串"META-SPL"的字节数组byte[] out.write(JobSplit.META_SPLIT_FILE_HEADER); // 写入分片元数据版本号splitMetaInfoVersion,当前为1 WritableUtils.writeVInt(out, splitMetaInfoVersion); // 写入分片元数据个数,为分片元数据信息SplitMetaInfo数组个数allSplitMetaInfo.length WritableUtils.writeVInt(out, allSplitMetaInfo.length); // 遍历分片元数据信息SplitMetaInfo数组allSplitMetaInfo中每个splitMetaInfo,挨个写入输出流 for (JobSplit.SplitMetaInfo splitMetaInfo : allSplitMetaInfo) { splitMetaInfo.write(out); } // 关闭输出流out out.close(); } writeJobSplitMetaInfo()方法的主体逻辑也十分清晰,大体如下: 1、调用HDFS文件系统FileSystem的create()方法,生成分片元数据信息文件,并获取文件系统数据输出流FSDataOutputStream实例out,对应文件路径为jobSubmitDir/job.splitmetainfo,jobSubmitDir为参数yarn.app.mapreduce.am.staging-dir指定的路径/作业所属用户user/.staging/作业ID,对应权限为JobSubmissionFiles.JOB_FILE_PERMISSION,即0644,rw-r--r--; 2、写入分片元数据头部信息UTF-8格式的字符串"META-SPL"的字节数组byte[]; 3、写入分片元数据版本号splitMetaInfoVersion,当前为1; 4、写入分片元数据个数,为分片元数据信息SplitMetaInfo数组个数allSplitMetaInfo.length; 5、遍历分片元数据信息SplitMetaInfo数组allSplitMetaInfo中每个splitMetaInfo,挨个写入输出流; 6、关闭输出流out。 我们看下如何序列化JobSplit.SplitMetaInfo,将其写入文件,JobSplit.SplitMetaInfo的write()如下: public void write(DataOutput out) throws IOException { // 将分片位置个数写入分片元数据信息文件 WritableUtils.writeVInt(out, locations.length); // 遍历位置信息,写入分片元数据信息文件 for (int i = 0; i < locations.length; i++) { Text.writeString(out, locations[i]); } // 写入分片元数据信息的起始位置 WritableUtils.writeVLong(out, startOffset); // 写入分片大小 WritableUtils.writeVLong(out, inputDataLength); } 每个分片的元数据信息,包括分片位置个数、分片文件位置、分片元数据信息的起始位置、分片大小等内容。 总结 JobSplitWriter被作业客户端用于写分片相关文件,包括分片数据文件job.split和分片元数据信息文件job.splitmetainfo。分片数据文件job.split存储的主要是每个分片对应的HDFS文件路径,和其在HDFS文件中的起始位置、长度等信息,而分片元数据信息文件job.splitmetainfo存储的则是每个分片在分片数据文件job.split中的起始位置、分片大小等信息。 job.split文件内容:文件头 + 分片 + 分片 + ... + 分片 文件头:"SPL" + 版本号1 分片:分片类 + 分片数据,分片类=String类型split对应类名,分片数据=String类型HDFS文件路径全名+Long类型分片在HDFS文件中的起始位置+Long类型分片在HDFS文件中的长度 job.splitmetainfo文件内容:文件头 + 分片元数据个数 +分片元数据 +分片元数据 + ... +分片元数据 文件头:"META-SPL" + 版本号1 分片元数据个数:分片元数据的个数 分片元数据:分片位置个数+分片位置+在分片文件job.split中的起始位置+分片大小

优秀的个人博客,低调大师

MapReduce源码分析之InputFormat

InputFormat描述了一个Map-Reduce作业中的输入规范。Map-Reduce框架依靠作业的InputFormat实现以下内容: 1、校验作业的输入规范; 2、分割输入文件(可能为多个),生成逻辑输入分片InputSplit(往往为多个),每个输入分片InputSplit接着被分配给单独的Mapper; 3、提供记录读取器RecordReader的实现,RecordReader被用于从逻辑输入分片InputSplit收集输入记录,这些输入记录会被交由Mapper处理。 基于文件的输入格式的默认行为,作为代表性的子类FileInputFormat,基于输入文件的总大小(单位byte)来切分成逻辑输入分片InputSplit。然而,输入文件的文件系统数据块大小,被用作输入分片大小的上界。输入分片大小的下界则可以在mapred-default.xml配置文件中通过参数mapreduce.input.fileinputformat.split.minsize来配置。 无疑,由于记录界限应该被遵守,基于输入大小的逻辑输入分片不满足很多应用。在这种情况下,应用不得不实现一个记录阅读器RecordReader,以便遵守记录边界,并提出一个面向记录的逻辑输入分片视图给单个任务。 InputFormat是一个抽象类,其中,实现分片的是getSplits()方法,其定义如下: public abstract List<InputSplit> getSplits(JobContext context ) throws IOException, InterruptedException; getSplits()方法为作业在逻辑上切分输入文件集合 。每个输入分片将会被分配给单个Mapper进行处理。注意,这个切分只是对输入进行逻辑上的切分,输入文件并不会在物理上被分割成块。比如,一个分片可能是<输入文件路径,起始位置,长度>元组。InputFormat也会创建记录阅读器RecordReader去读取这个输入分片InputSplit。 而提供记录阅读器的是createRecordReader()方法,其定义如下: public abstract RecordReader<K,V> createRecordReader(InputSplit split, TaskAttemptContext context ) throws IOException, InterruptedException; createRecordReader()方法为给定分片创建一个记录阅读器。在分片被使用之前,框架将调用RecordReader的initialize(InputSplit, TaskAttemptContext)方法完成初始化。它需要两个参数: 1、InputSplit split:需要被读入的分片; 2、TaskAttemptContext context:任务上下文,存储了任务的相关信息。

资源下载

更多资源
Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

用户登录
用户注册