首页 文章 精选 留言 我的

精选列表

搜索[storm],共1270篇文章
优秀的个人博客,低调大师

Storm-源码分析- timer (backtype.storm.timer)

mk-timer timer是基于PriorityQueue实现的(和PriorityBlockingQueue区别, 在于没有阻塞机制, 不是线程安全的), 优先级队列是堆数据结构的典型应用 默认情况下, 按照自然顺序(其实就是默认comparator的定义), 最小的元素排在堆头 当然也可以自己重新实现comparator接口, 比如timer就用reify重新实现了comparator接口 整个过程其实比较简单, 开个timer-thread, 不断check PriorityQueue里面时间最小的timer是否已经可以触发 如果可以, 就poll出来, 调用callback, 并sleep, 都很好理解 唯一需要说的是, 这里使用Semaphore, 信号量和lock相似, 都是用于互斥 不同在于, 信号量模拟资源管理, 所以不同于lock的排他, 信号量可以接收多个aquire(取决于配置) 另外一个比较大的区别, lock是解铃还须系铃人, 谁锁谁解, 而信号量无所谓, 任何线程都可以调用release, 或acquire 这里使用信号量, 是用于在cancel-timer时, 等待timer-thread结束 (defn cancel-timer [timer] (check-active! timer) (locking (:lock timer) (reset! (:active timer) false) (.interrupt (:timer-thread timer))) (.acquire (:cancel-notifier timer))) 因为cancel的过程就是将active置false, 然后就是调用acquire等待信号量cancel-notifier被释放 而timer-thread在线程结束前, 会release这个信号量 (defnk mk-timer [:kill-fn (fn [& _] )] (let [queue (PriorityQueue. 10 (reify Comparator (compare [this o1 o2] (- (first o1) (first o2)) ) (equals [this obj] true ))) active (atom true) ;;标志位 lock (Object.) ;;创建lock对象, 由于PriorityQueue非线程安全, 所以使用locking来保证同时只有一个线程访问queue notifier (Semaphore. 0) ;;创建信号量, 初始为0 timer-thread (Thread. (fn [] (while @active (try ;;peek读但不从queue中取出, 先读出time看看, 符合条件再取出 (let [[time-secs _ _ :as elem] (locking lock (.peek queue))] (if (and elem (>= (current-time-secs) time-secs)) ;;无法保证恰好, 只要当前时间>=time-secs, 就可以执行, 可想而知对于afn必须不能耗时, 否则会影响其他timer ;; imperative to not run the function inside the timer lock ;; otherwise, it's possible to deadlock if function deals with other locks ;; (like the submit lock) (let [afn (locking lock (second (.poll queue)))] ;;poll从queue中取出 (afn)) ;;真正执行timer中的callback (Time/sleep 1000) )) (catch Throwable t ;; because the interrupted exception can be wrapped in a runtimeexception (when-not (exception-cause? InterruptedException t) (kill-fn t) (reset! active false) (throw t)) ))) (.release notifier)))] (.setDaemon timer-thread true) (.setPriority timer-thread Thread/MAX_PRIORITY) (.start timer-thread) {:timer-thread timer-thread :queue queue :active active :lock lock :cancel-notifier notifier})) schedule schedule其实就是往PriorityQueue里面插入timer 对于循环schdule, 就是在timer的callback里面, 再次schedule (defnk schedule [timer delay-secs afn :check-active true] (when check-active (check-active! timer)) (let [id (uuid) ^PriorityQueue queue (:queue timer)] (locking (:lock timer) (.add queue [(+ (current-time-secs) delay-secs) afn id]) ))) (defn schedule-recurring [timer delay-secs recur-secs afn] (schedule timer delay-secs (fn this [] (afn) (schedule timer recur-secs this :check-active false)) ; this avoids a race condition with cancel-timer )) 使用例子 Supervisor中的使用例子, 定期的调用hb函数更新supervisor的hb 在mk-timer时, 传入的kill-fn callback, 会在timer-thread发生exception的时候被调用 :timer (mk-timer :kill-fn (fn [t] (log-error t "Error when processing event") (halt-process! 20 "Error when processing an event") )) (schedule-recurring (:timer supervisor) 0 (conf SUPERVISOR-HEARTBEAT-FREQUENCY-SECS) heartbeat-fn) 本文章摘自博客园,原文发布日期:2013-07-02

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

Storm-源码分析-LocalState (backtype.storm.utils)

LocalState A simple, durable, atomic K/V database. *Very inefficient*, should only be used for occasional reads/writes. Every read/write hits disk. 基于map实现, 每次读写都需要从磁盘上将数据读出, 并反序列化成map, 这个过程称为snapshot. 所以说是比较简单和低效的, 只能用于读取配置或参数, 这种偶尔读取的场景. public synchronized Map<Object, Object> snapshot() throws IOException { int attempts = 0; while(true) { String latestPath = _vs.mostRecentVersionPath(); if(latestPath==null) return new HashMap<Object, Object>(); try { return (Map<Object, Object>) Utils.deserialize(FileUtils.readFileToByteArray(new File(latestPath))); } catch(IOException e) { attempts++; if(attempts >= 10) { throw e; } } } 读写操作都是基于map的操作, get和put, 但是put需要做persist操作. 这里使用synchronized来做对象的线程间同步, 对于一个LocalState对象, 所有synchronized标有的函数只能被串行操作. public Object get(Object key) throws IOException { return snapshot().get(key); } public synchronized void put(Object key, Object val, boolean cleanup) throws IOException { Map<Object, Object> curr = snapshot(); curr.put(key, val); persist(curr, cleanup); } 当然不止这么简单, 为了达到atomic, 还使用了VersionedStore, 参考下一章 persist不会去update现有的文件, 而是不断的产生递增version的文件, 故每一批更新都会产生一个新的文件 把需要写入的数据序列化 创建新的versionfile的path 把数据写入versionfile 调用succeedVersion, 创建tokenfile以标志versionfile的写入完成 清除旧版本, 只保留4个版本 private void persist(Map<Object, Object> val, boolean cleanup) throws IOException { byte[] toWrite = Utils.serialize(val); String newPath = _vs.createVersion(); FileUtils.writeByteArrayToFile(new File(newPath), toWrite); _vs.succeedVersion(newPath); if(cleanup) _vs.cleanup(4); } VersionedStore public VersionedStore(String path) throws IOException { _root = path; mkdirs(_root); } 这个store, 其实就是_root目录下的一堆文件 文件分两种, VersionFile, _root + version, 真正的数据存储文件 TokenFile, _root + version + “.version”, 标志位文件, 标志version文件是否完成写操作, 以避免读到正在更新的文件 getAllVersions就是读出所有_root目录下的所有完成写操作的文件, 读出version, 并做从大到小的排序 public List<Long> getAllVersions() throws IOException { List<Long> ret = new ArrayList<Long>(); for(String s: listDir(_root)) { if(s.endsWith(FINISHED_VERSION_SUFFIX)) { ret.add(validateAndGetVersion(s)); } } Collections.sort(ret); Collections.reverse(ret); return ret; } 找到最新的版本文件 public Long mostRecentVersion() throws IOException { List<Long> all = getAllVersions(); if(all.size()==0) return null; return all.get(0); 创建新版本号, 用当前时间作为version public String createVersion() throws IOException { Long mostRecent = mostRecentVersion(); long version = Time.currentTimeMillis(); if(mostRecent!=null && version <= mostRecent) { version = mostRecent + 1; } return createVersion(version); } public String createVersion(long version) throws IOException { String ret = versionPath(version); if(getAllVersions().contains(version)) throw new RuntimeException("Version already exists or data already exists"); else return ret; } 创建tokenfile, 以标记versionfile写完成 public void succeedVersion(String path) throws IOException { long version = validateAndGetVersion(path); // should rewrite this to do a file move createNewFile(tokenPath(version)); } 清除旧的版本, 只保留versionsToKeep个, 清除操作就是删除versionfile和tokenfile public void cleanup(int versionsToKeep) throws IOException { List<Long> versions = getAllVersions(); if(versionsToKeep >= 0) { versions = versions.subList(0, Math.min(versions.size(), versionsToKeep)); } HashSet<Long> keepers = new HashSet<Long>(versions); for(String p: listDir(_root)) { Long v = parseVersion(p); if(v!=null && !keepers.contains(v)) { deleteVersion(v); } } } 本文章摘自博客园,原文发布日期:2013-07-11

资源下载

更多资源
腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Nacos

Nacos

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

Sublime Text

Sublime Text

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

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册