• 设为首页
  • 点击收藏
  • 手机版
    手机扫一扫访问
    迪恩网络手机版
  • 关注官方公众号
    微信扫一扫关注
    公众号

Java TimeCacheMap类代码示例

原作者: [db:作者] 来自: [db:来源] 收藏 邀请

本文整理汇总了Java中backtype.storm.utils.TimeCacheMap的典型用法代码示例。如果您正苦于以下问题:Java TimeCacheMap类的具体用法?Java TimeCacheMap怎么用?Java TimeCacheMap使用的例子?那么恭喜您, 这里精选的类代码示例或许可以为您提供帮助。



TimeCacheMap类属于backtype.storm.utils包,在下文中一共展示了TimeCacheMap类的20个代码示例,这些例子默认根据受欢迎程度排序。您可以为喜欢或者感觉有用的代码点赞,您的评价将有助于我们的系统推荐出更棒的Java代码示例。

示例1: finishFileUpload

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void finishFileUpload(String location) throws TException {

	TimeCacheMap<Object, Object> uploaders = data.getUploaders();
	Object obj = uploaders.get(location);
	if (obj == null) {
		throw new TException(
				"File for that location does not exist (or timed out)");
	}
	try {
		if (obj instanceof WritableByteChannel) {
			WritableByteChannel channel = (WritableByteChannel) obj;
			channel.close();
			uploaders.remove(location);
			LOG.info("Finished uploading file from client: " + location);
		} else {
			throw new TException("Object isn't WritableByteChannel for "
					+ location);
		}
	} catch (IOException e) {
		LOG.error(" WritableByteChannel close failed when finishFileUpload "
				+ location);
	}

}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:26,代码来源:ServiceHandler.java


示例2: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void prepare(Map conf, TopologyContext context,
		OutputCollector collector) {

	taskId = String.valueOf(context.getThisTaskId());
	taskName = context.getThisComponentId() + "_" + context.getThisTaskId();

	this.basicCollector = new BasicOutputCollector(collector);
	this.collector = collector;

	if (delegate instanceof ICommitter) {
		isCommiter = true;
		commited = new TimeCacheMap<Object, Object>(
				context.maxTopologyMessageTimeout());
		mkCommitDir(conf);
	}

	delegate.prepare(conf, context);

}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:21,代码来源:CoordinatedBolt.java


示例3: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
public void prepare(Map conf, TopologyContext context,
		OutputCollector collector) {

	taskId = String.valueOf(context.getThisTaskId());
	taskName = context.getThisComponentId() + "_" + context.getThisTaskId();

	this.basicCollector = new BasicOutputCollector(collector);
	this.collector = collector;

	if (delegate instanceof ICommitter) {
		isCommiter = true;
		commited = new TimeCacheMap<Object, Object>(
				context.maxTopologyMessageTimeout());
		mkCommitDir(conf);
	}

	delegate.prepare(conf, context);

}
 
开发者ID:songtk,项目名称:learn_jstorm,代码行数:20,代码来源:CoordinatedBolt.java


示例4: createFileHandler

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
public void createFileHandler() {
    TimeCacheMap.ExpiredCallback<Object, Object> expiredCallback = new TimeCacheMap.ExpiredCallback<Object, Object>() {
        @Override
        public void expire(Object key, Object val) {
            try {
                LOG.info("Close file " + String.valueOf(key));
                if (val != null) {
                    if (val instanceof Channel) {
                        Channel channel = (Channel) val;
                        channel.close();
                    } else if (val instanceof BufferFileInputStream) {
                        BufferFileInputStream is = (BufferFileInputStream) val;
                        is.close();
                    }
                }
            } catch (IOException e) {
                LOG.error(e.getMessage(), e);
            }

        }
    };

    int file_copy_expiration_secs = JStormUtils.parseInt(conf.get(Config.NIMBUS_FILE_COPY_EXPIRATION_SECS), 30);
    uploaders = new TimeCacheMap<Object, Object>(file_copy_expiration_secs, expiredCallback);
    downloaders = new TimeCacheMap<Object, Object>(file_copy_expiration_secs, expiredCallback);
}
 
开发者ID:kkllwww007,项目名称:jstrom,代码行数:27,代码来源:NimbusData.java


示例5: uploadChunk

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
/**
 * uploading topology jar data
 */
@Override
public void uploadChunk(String location, ByteBuffer chunk) throws TException {
    TimeCacheMap<Object, Object> uploaders = data.getUploaders();
    Object obj = uploaders.get(location);
    if (obj == null) {
        throw new TException("File for that location does not exist (or timed out) " + location);
    }
    try {
        if (obj instanceof WritableByteChannel) {
            WritableByteChannel channel = (WritableByteChannel) obj;
            channel.write(chunk);
            uploaders.put(location, channel);
        } else {
            throw new TException("Object isn't WritableByteChannel for " + location);
        }
    } catch (IOException e) {
        String errMsg = " WritableByteChannel write filed when uploadChunk " + location;
        LOG.error(errMsg);
        throw new TException(e);
    }

}
 
开发者ID:kkllwww007,项目名称:jstrom,代码行数:26,代码来源:ServiceHandler.java


示例6: finishFileUpload

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void finishFileUpload(String location) throws TException {
    TimeCacheMap<Object, Object> uploaders = data.getUploaders();
    Object obj = uploaders.get(location);
    if (obj == null) {
        throw new TException("File for that location does not exist (or timed out)");
    }
    try {
        if (obj instanceof WritableByteChannel) {
            WritableByteChannel channel = (WritableByteChannel) obj;
            channel.close();
            uploaders.remove(location);
            LOG.info("Finished uploading file from client: " + location);
        } else {
            throw new TException("Object isn't WritableByteChannel for " + location);
        }
    } catch (IOException e) {
        LOG.error(" WritableByteChannel close failed when finishFileUpload " + location);
    }

}
 
开发者ID:kkllwww007,项目名称:jstrom,代码行数:22,代码来源:ServiceHandler.java


示例7: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
public void prepare(Map conf, TopologyContext context, OutputCollector collector) {

        taskId = String.valueOf(context.getThisTaskId());
        taskName = context.getThisComponentId() + "_" + context.getThisTaskId();

        this.basicCollector = new BasicOutputCollector(collector);
        this.collector = collector;

        if (delegate instanceof ICommitter) {
            isCommiter = true;
            commited = new TimeCacheMap<Object, Object>(context.maxTopologyMessageTimeout());
            mkCommitDir(conf);
        }

        delegate.prepare(conf, context);

    }
 
开发者ID:kkllwww007,项目名称:jstrom,代码行数:18,代码来源:CoordinatedBolt.java


示例8: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
public void prepare(Map config, TopologyContext context, OutputCollector collector) {
    TimeCacheMap.ExpiredCallback<Object, TrackingInfo> callback = null;
    if (_delegate instanceof TimeoutCallback) {
        callback = new TimeoutItems();
    }
    _tracked = new TimeCacheMap<Object, TrackingInfo>(context.maxTopologyMessageTimeout(), callback);
    _collector = collector;
    _delegate.prepare(config, context, new OutputCollector(new CoordinatedOutputCollector(collector)));
    for (String component : Utils.get(context.getThisTargets(), Constants.COORDINATED_STREAM_ID, new HashMap<String, Grouping>()).keySet()) {
        for (Integer task : context.getComponentTasks(component)) {
            _countOutTasks.add(task);
        }
    }
    if (!_sourceArgs.isEmpty()) {
        _numSourceReports = 0;
        for (Entry<String, SourceArgs> entry : _sourceArgs.entrySet()) {
            if (entry.getValue().singleCount) {
                _numSourceReports += 1;
            } else {
                _numSourceReports += context.getComponentTasks(entry.getKey()).size();
            }
        }
    }
}
 
开发者ID:kkllwww007,项目名称:jstrom,代码行数:25,代码来源:CoordinatedBolt.java


示例9: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
public void prepare(Map config, TopologyContext context, OutputCollector collector) {
    TimeCacheMap.ExpiredCallback<Object, TrackingInfo> callback = null;
    if (_delegate instanceof TimeoutCallback) {
        callback = new TimeoutItems();
    }
    _tracked = new TimeCacheMap<>(context.maxTopologyMessageTimeout(), callback);
    _collector = collector;
    _delegate.prepare(config, context, new OutputCollector(new CoordinatedOutputCollector(collector)));
    for (String component : Utils.get(context.getThisTargets(), Constants.COORDINATED_STREAM_ID, new HashMap<String, Grouping>()).keySet()) {
        for (Integer task : context.getComponentTasks(component)) {
            _countOutTasks.add(task);
        }
    }
    if (!_sourceArgs.isEmpty()) {
        _numSourceReports = 0;
        for (Entry<String, SourceArgs> entry : _sourceArgs.entrySet()) {
            if (entry.getValue().singleCount) {
                _numSourceReports += 1;
            } else {
                _numSourceReports += context.getComponentTasks(entry.getKey()).size();
            }
        }
    }
}
 
开发者ID:alibaba,项目名称:jstorm,代码行数:25,代码来源:CoordinatedBolt.java


示例10: createNimbusData

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
private NimbusData createNimbusData(Map conf, INimbus inimbus)
		throws Exception {

	TimeCacheMap.ExpiredCallback<Object, Object> expiredCallback = new TimeCacheMap.ExpiredCallback<Object, Object>() {
		@Override
		public void expire(Object key, Object val) {
			try {
				LOG.info("Close file " + String.valueOf(key));
				if (val != null) {
					if (val instanceof Channel) {
						Channel channel = (Channel) val;
						channel.close();
					} else if (val instanceof BufferFileInputStream) {
						BufferFileInputStream is = (BufferFileInputStream) val;
						is.close();
					}
				}
			} catch (IOException e) {
				LOG.error(e.getMessage(), e);
			}

		}
	};

	int file_copy_expiration_secs = JStormUtils.parseInt(
			conf.get(Config.NIMBUS_FILE_COPY_EXPIRATION_SECS), 30);
	TimeCacheMap<Object, Object> uploaders = new TimeCacheMap<Object, Object>(
			file_copy_expiration_secs, expiredCallback);
	TimeCacheMap<Object, Object> downloaders = new TimeCacheMap<Object, Object>(
			file_copy_expiration_secs, expiredCallback);

	// Callback callback=new TimerCallBack();
	// StormTimer timer=Timer.mkTimerTimer(callback);
	NimbusData data = new NimbusData(conf, downloaders, uploaders, inimbus);

	return data;

}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:39,代码来源:NimbusServer.java


示例11: uploadChunk

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
/**
 * uploading topology jar data
 */
@Override
public void uploadChunk(String location, ByteBuffer chunk)
		throws TException {
	TimeCacheMap<Object, Object> uploaders = data.getUploaders();
	Object obj = uploaders.get(location);
	if (obj == null) {
		throw new TException(
				"File for that location does not exist (or timed out) "
						+ location);
	}
	try {
		if (obj instanceof WritableByteChannel) {
			WritableByteChannel channel = (WritableByteChannel) obj;
			channel.write(chunk);
			uploaders.put(location, channel);
		} else {
			throw new TException("Object isn't WritableByteChannel for "
					+ location);
		}
	} catch (IOException e) {
		String errMsg = " WritableByteChannel write filed when uploadChunk "
				+ location;
		LOG.error(errMsg);
		throw new TException(e);
	}

}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:31,代码来源:ServiceHandler.java


示例12: downloadChunk

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public ByteBuffer downloadChunk(String id) throws TException {
	TimeCacheMap<Object, Object> downloaders = data.getDownloaders();
	Object obj = downloaders.get(id);
	if (obj == null) {
		throw new TException("Could not find input stream for that id");
	}

	try {
		if (obj instanceof BufferFileInputStream) {

			BufferFileInputStream is = (BufferFileInputStream) obj;
			byte[] ret = is.read();
			if (ret != null) {
				downloaders.put(id, is);
				return ByteBuffer.wrap(ret);
			}
		} else {
			throw new TException("Object isn't BufferFileInputStream for "
					+ id);
		}
	} catch (IOException e) {
		LOG.error("BufferFileInputStream read failed when downloadChunk ",
				e);
		throw new TException(e);
	}
	byte[] empty = {};
	return ByteBuffer.wrap(empty);
}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:30,代码来源:ServiceHandler.java


示例13: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void prepare(Map conf, TopologyContext context,
		OutputCollector collector) {
	_fieldLocations = new HashMap<String, GlobalStreamId>();
	_collector = collector;
	int timeout = ((Number) conf.get(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS))
			.intValue();
	_pending = new TimeCacheMap<List<Object>, Map<GlobalStreamId, Tuple>>(
			timeout, new ExpireCallback());
	_numSources = context.getThisSources().size();
	Set<String> idFields = null;
	for (GlobalStreamId source : context.getThisSources().keySet()) {
		Fields fields = context.getComponentOutputFields(
				source.get_componentId(), source.get_streamId());
		Set<String> setFields = new HashSet<String>(fields.toList());
		if (idFields == null)
			idFields = setFields;
		else
			idFields.retainAll(setFields);

		for (String outfield : _outFields) {
			for (String sourcefield : fields) {
				if (outfield.equals(sourcefield)) {
					_fieldLocations.put(outfield, source);
				}
			}
		}
	}
	_idFields = new Fields(new ArrayList<String>(idFields));

	if (_fieldLocations.size() != _outFields.size()) {
		throw new RuntimeException(
				"Cannot find all outfields among sources");
	}
}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:36,代码来源:SingleJoinBolt.java


示例14: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void prepare(Map config, TopologyContext context,
		OutputCollector collector) {
	TimeCacheMap.ExpiredCallback<Object, TrackingInfo> callback = null;
	if (_delegate instanceof TimeoutCallback) {
		callback = new TimeoutItems();
	}
	_tracked = new TimeCacheMap<Object, TrackingInfo>(
			context.maxTopologyMessageTimeout(), callback);
	_collector = collector;
	_delegate.prepare(config, context, new OutputCollector(
			new CoordinatedOutputCollector(collector)));
	for (String component : Utils.get(context.getThisTargets(),
			Constants.COORDINATED_STREAM_ID,
			new HashMap<String, Grouping>()).keySet()) {
		for (Integer task : context.getComponentTasks(component)) {
			_countOutTasks.add(task);
		}
	}
	if (!_sourceArgs.isEmpty()) {
		_numSourceReports = 0;
		for (Entry<String, SourceArgs> entry : _sourceArgs.entrySet()) {
			if (entry.getValue().singleCount) {
				_numSourceReports += 1;
			} else {
				_numSourceReports += context.getComponentTasks(
						entry.getKey()).size();
			}
		}
	}
}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:32,代码来源:CoordinatedBolt.java


示例15: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void prepare(Map stormConf, TopologyContext context) {
	this.conf = stormConf;

	int timeoutSeconds = JStormUtils.parseInt(
			conf.get(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS), 30);
	counters = new TimeCacheMap<BatchId, AtomicLong>(timeoutSeconds);
	
	LOG.info("Successfully do prepare");

}
 
开发者ID:zhangjunfang,项目名称:jstorm-0.9.6.3-,代码行数:12,代码来源:CountBolt.java


示例16: createNimbusData

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
/**
 * 1.此处定义一个TimcacheMap缓存失效时候的callback函数,关闭channel和stream
 * 建立上传和下载两个缓存,缓存流
 * 新建一个Nimbus对象
 */
@SuppressWarnings("rawtypes")
private NimbusData createNimbusData(Map conf, INimbus inimbus) throws Exception {

    TimeCacheMap.ExpiredCallback<Object, Object> expiredCallback = new TimeCacheMap.ExpiredCallback<Object, Object>() {
        @Override
        public void expire(Object key, Object val) {
            try {
                LOG.info("Close file " + String.valueOf(key));
                if (val != null) {
                    if (val instanceof Channel) {
                        Channel channel = (Channel) val;
                        channel.close();
                    } else if (val instanceof BufferFileInputStream) {
                        BufferFileInputStream is = (BufferFileInputStream) val;
                        is.close();
                    }
                }
            } catch (IOException e) {
                LOG.error(e.getMessage(), e);
            }

        }
    };

    int file_copy_expiration_secs = JStormUtils.parseInt(conf.get(Config.NIMBUS_FILE_COPY_EXPIRATION_SECS), 30);
    TimeCacheMap<Object, Object> uploaders = new TimeCacheMap<Object, Object>(file_copy_expiration_secs,
            expiredCallback);
    TimeCacheMap<Object, Object> downloaders = new TimeCacheMap<Object, Object>(file_copy_expiration_secs,
            expiredCallback);

    // Callback callback=new TimerCallBack();
    // StormTimer timer=Timer.mkTimerTimer(callback);
    NimbusData data = new NimbusData(conf, downloaders, uploaders, inimbus);

    return data;

}
 
开发者ID:songtk,项目名称:learn_jstorm,代码行数:43,代码来源:NimbusServer.java


示例17: downloadChunk

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public ByteBuffer downloadChunk(String id) throws TException {
	TimeCacheMap<Object, Object> downloaders = data.getDownloaders();
	Object obj = downloaders.get(id);
	if (obj == null) {
		throw new TException("Could not find input stream for that id");
	}

	try {
		if (obj instanceof BufferFileInputStream) {

			BufferFileInputStream is = (BufferFileInputStream) obj;
			byte[] ret = is.read();
			if (ret != null) {
				downloaders.put(id, (BufferFileInputStream) is);
				return ByteBuffer.wrap(ret);
			}
		} else {
			throw new TException("Object isn't BufferFileInputStream for "
					+ id);
		}
	} catch (IOException e) {
		LOG.error("BufferFileInputStream read failed when downloadChunk ",
				e);
		throw new TException(e);
	}
	byte[] empty = {};
	return ByteBuffer.wrap(empty);
}
 
开发者ID:songtk,项目名称:learn_jstorm,代码行数:30,代码来源:ServiceHandler.java


示例18: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
public void prepare(Map config, TopologyContext context,
		OutputCollector collector) {
	TimeCacheMap.ExpiredCallback<Object, TrackingInfo> callback = null;
	if (_delegate instanceof TimeoutCallback) {
		callback = new TimeoutItems();
	}
	_tracked = new TimeCacheMap<Object, TrackingInfo>(
			context.maxTopologyMessageTimeout(), callback);
	_collector = collector;
	_delegate.prepare(config, context, new OutputCollector(
			new CoordinatedOutputCollector(collector)));
	for (String component : Utils.get(context.getThisTargets(),
			Constants.COORDINATED_STREAM_ID,
			new HashMap<String, Grouping>()).keySet()) {
		for (Integer task : context.getComponentTasks(component)) {
			_countOutTasks.add(task);
		}
	}
	if (!_sourceArgs.isEmpty()) {
		_numSourceReports = 0;
		for (Entry<String, SourceArgs> entry : _sourceArgs.entrySet()) {
			if (entry.getValue().singleCount) {
				_numSourceReports += 1;
			} else {
				_numSourceReports += context.getComponentTasks(
						entry.getKey()).size();
			}
		}
	}
}
 
开发者ID:songtk,项目名称:learn_jstorm,代码行数:31,代码来源:CoordinatedBolt.java


示例19: prepare

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public void prepare(Map conf, TopologyContext context, OutputCollector collector) {
  _fieldLocations = new HashMap<String, GlobalStreamId>();
  _collector = collector;
  int timeout = ((Number) conf.get(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS)).intValue();
  _pending = new TimeCacheMap<List<Object>, Map<GlobalStreamId, Tuple>>(timeout, new ExpireCallback());
  _numSources = context.getThisSources().size();
  Set<String> idFields = null;
  for (GlobalStreamId source : context.getThisSources().keySet()) {
    Fields fields = context.getComponentOutputFields(source.get_componentId(), source.get_streamId());
    Set<String> setFields = new HashSet<String>(fields.toList());
    if (idFields == null)
      idFields = setFields;
    else
      idFields.retainAll(setFields);

    for (String outfield : _outFields) {
      for (String sourcefield : fields) {
        if (outfield.equals(sourcefield)) {
          _fieldLocations.put(outfield, source);
        }
      }
    }
  }
  _idFields = new Fields(new ArrayList<String>(idFields));

  if (_fieldLocations.size() != _outFields.size()) {
    throw new RuntimeException("Cannot find all outfields among sources");
  }
}
 
开发者ID:luozhaoyu,项目名称:big-data-system,代码行数:31,代码来源:SingleJoinBolt.java


示例20: downloadChunk

import backtype.storm.utils.TimeCacheMap; //导入依赖的package包/类
@Override
public ByteBuffer downloadChunk(String id) throws TException {
    TimeCacheMap<Object, Object> downloaders = data.getDownloaders();
    Object obj = downloaders.get(id);
    if (obj == null) {
        throw new TException("Could not find input stream for that id");
    }

    try {
        if (obj instanceof BufferFileInputStream) {

            BufferFileInputStream is = (BufferFileInputStream) obj;
            byte[] ret = is.read();
            if (ret != null) {
                downloaders.put(id, is);
                return ByteBuffer.wrap(ret);
            }
        } else {
            throw new TException("Object isn't BufferFileInputStream for " + id);
        }
    } catch (IOException e) {
        LOG.error("BufferFileInputStream read failed when downloadChunk ", e);
        throw new TException(e);
    }
    byte[] empty = {};
    return ByteBuffer.wrap(empty);
}
 
开发者ID:kkllwww007,项目名称:jstrom,代码行数:28,代码来源:ServiceHandler.java



注:本文中的backtype.storm.utils.TimeCacheMap类示例整理自Github/MSDocs等源码及文档管理平台,相关代码片段筛选自各路编程大神贡献的开源项目,源码版权归原作者所有,传播和使用请参考对应项目的License;未经允许,请勿转载。


鲜花

握手

雷人

路过

鸡蛋
该文章已有0人参与评论

请发表评论

全部评论

专题导读
上一篇:
Java RendererCapabilities类代码示例发布时间:2022-05-21
下一篇:
Java EnumMovingObjectType类代码示例发布时间:2022-05-21
热门推荐
阅读排行榜

扫描微信二维码

查看手机版网站

随时了解更新最新资讯

139-2527-9053

在线客服(服务时间 9:00~18:00)

在线QQ客服
地址:深圳市南山区西丽大学城创智工业园
电邮:jeky_zhao#qq.com
移动电话:139-2527-9053

Powered by 互联科技 X3.4© 2001-2213 极客世界.|Sitemap