Во-первых, вам нужно создать класс событий, который будет определять все атрибуты, которые имеет ваш поток событий, а затем создавать все методы получения, установки и другие методы. Примером этого класса будет
public class qrsIntervalStreamEvent {
public Integer Sensor_id;
public long time;
public Integer qrsInterval;
public qrsIntervalStreamEvent(Integer sensor_id, long time, Integer qrsInterval) {
Sensor_id = sensor_id;
this.time = time;
this.qrsInterval = qrsInterval;
}
public Integer getSensor_id() {
return Sensor_id;
}
public void setSensor_id(Integer sensor_id) {
Sensor_id = sensor_id;
}
public long getTime() {
return time;
}
public void setTime(long time) {
this.time = time;
}
public Integer getQrsInterval() {
return qrsInterval;
}
public void setQrsInterval(Integer qrsInterval) {
this.qrsInterval = qrsInterval;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof qrsIntervalStreamEvent)) return false;
qrsIntervalStreamEvent that = (qrsIntervalStreamEvent) o;
if (getTime() != that.getTime()) return false;
if (getSensor_id() != null ? !getSensor_id().equals(that.getSensor_id()) : that.getSensor_id() != null)
return false;
return getQrsInterval() != null ? getQrsInterval().equals(that.getQrsInterval()) : that.getQrsInterval() == null;
}
@Override
public int hashCode() {
int result = getSensor_id() != null ? getSensor_id().hashCode() : 0;
result = 31 * result + (int) (getTime() ^ (getTime() >>> 32));
result = 31 * result + (getQrsInterval() != null ? getQrsInterval().hashCode() : 0);
return result;
}
@Override
public String toString() {
return "StreamEvent{" +
"Sensor_id=" + Sensor_id +
", time=" + time +
", qrsInterval=" + qrsInterval +
'}';
}
} //class
Теперь предположим, что вы хотите отправить эти события через x событий / 5 секунд, тогда вы можете написать код примерно так
public class Qrs_interval_Gen extends RichParallelSourceFunction<qrsIntervalStreamEvent> {
@Override
public void run(SourceContext<qrsIntervalStreamEvent> sourceContext) throws Exception {
int qrsInterval;
int Sensor_id;
long currentTime;
Random random = new Random();
Integer InputRate = 10;
Integer Sleeptime = 1000 * 5 / InputRate ;
for(int i = 0 ; i <= 100000 ; i++){
// int randomNum = rand.nextInt((max - min) + 1) + min;
Sensor_id = 1;
qrsInterval = 10 + random.nextInt((20-10)+ 1);
// currentTime = System.currentTimeMillis();
currentTime = i;
//System.out.println("qrsInterval = " + qrsInterval + ", Sensor_id = "+ Sensor_id );
try {
Thread.sleep(Sleeptime);
} catch (InterruptedException e) {
e.printStackTrace();
}
qrsIntervalStreamEvent stream = new qrsIntervalStreamEvent(Sensor_id,currentTime,qrsInterval);
sourceContext.collect(stream);
} // for loop
}
@Override
public void cancel() {
}
}
Здесь вся логика осуществляется
если вы хотите отправлять x событий в секунду, время вашего сна будет обратным. Например, для отправки 10 событий в секунду
Время сна = 1000/10 = 100 миллисекунд
Точно так же для отправки 10 событий / 5 секунд время сна будет
Время сна = 1000 * 5/10 = 500 миллисекунд
Надеюсь, это поможет, дайте мне знать, если у вас возникнут вопросы
person
Amarjit Dhillon
schedule
29.10.2017