Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,10 @@ public class IncomingMessageEnvelope {
private final Object key;
private final Object message;
private final int size;
// the timestamp when this event occured, should be set by the event source, 0 means unassigned
// the System's nanoTime when this event occurred, should be set by the event source, 0 means unassigned
private long eventTime = 0L;
// the timestamp when this event is pickedup by samza, 0 means unassgined
private long arrivalTime = 0L;
// the System's nanoTime when this event is picked-up by samza, 0 means unassgined
private long arrivalTime;

/**
* Constructs a new IncomingMessageEnvelope from specified components.
Expand Down Expand Up @@ -70,7 +70,7 @@ public IncomingMessageEnvelope(SystemStreamPartition systemStreamPartition, Stri
this.key = key;
this.message = message;
this.size = size;
this.arrivalTime = Instant.now().toEpochMilli();
this.arrivalTime = System.nanoTime();
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,7 @@ public class SamzaSqlInputTransformer implements InputTransformer {
public Object apply(IncomingMessageEnvelope ime) {
Assert.notNull(ime, "ime is null");
KV<Object, Object> keyAndMessageKV = KV.of(ime.getKey(), ime.getMessage());
SamzaSqlRelMsgMetadata metadata = new SamzaSqlRelMsgMetadata(
(ime.getEventTime() == 0) ? "" : Instant.ofEpochMilli(ime.getEventTime()).toString(),
(ime.getArrivalTime() == 0) ? "" : Instant.ofEpochMilli(ime.getArrivalTime()).toString(), null);
SamzaSqlRelMsgMetadata metadata = new SamzaSqlRelMsgMetadata(ime.getEventTime(), ime.getArrivalTime(),0L);
SamzaSqlInputMessage samzaMsg = SamzaSqlInputMessage.of(keyAndMessageKV, metadata);
return samzaMsg;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ public SamzaSqlRelMessage convertToRelMessage(KV<Object, Object> samzaMessage) {
}

return new SamzaSqlRelMessage(samzaMessage.getKey(), payloadFieldNames, payloadFieldValues,
new SamzaSqlRelMsgMetadata("", "", ""));
new SamzaSqlRelMsgMetadata(0, 0, 0));
}

/**
Expand All @@ -108,7 +108,7 @@ public static SamzaSqlRelMessage convertToRelMessage(Object key, IndexedRecord r
List<Object> payloadFieldValues = new ArrayList<>();
fetchFieldNamesAndValuesFromIndexedRecord(record, payloadFieldNames, payloadFieldValues, schema);
return new SamzaSqlRelMessage(key, payloadFieldNames, payloadFieldValues,
new SamzaSqlRelMsgMetadata("", "", ""));
new SamzaSqlRelMsgMetadata(0, 0, 0));
}

public static void fetchFieldNamesAndValuesFromIndexedRecord(IndexedRecord record, List<String> fieldNames,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -130,15 +130,15 @@ public Object getKey() {
return key;
}

public void setEventTime(String eventTime) {
public void setEventTime(long eventTime) {
this.samzaSqlRelMsgMetadata.setEventTime(eventTime);
}

public void setArrivalTime(String arrivalTime) {
public void setArrivalTime(long arrivalTime) {
this.samzaSqlRelMsgMetadata.setArrivalTime(arrivalTime);
}

public void setScanTime(String scanTime) {
public void setScanTime(long scanTime) {
this.samzaSqlRelMsgMetadata.setScanTime(scanTime);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,60 +56,60 @@ public class SamzaSqlRelMsgMetadata implements Serializable {
* TODO: copy eventTime through from source to RelMessage
*/
@JsonProperty("eventTime")
private String eventTime;
private long eventTime;

/**
* the timestamp of when Samza App received the event
* TODO: set arrivalTime during conversion from IME to SamzaMessage
*/
@JsonProperty("arrivalTime")
private String arrivalTime;
private long arrivalTime;

/**
* the timestamp when SamzaSQL query starts processing the event
* set by the SamzaSQL Scan operator
*/
@JsonProperty("scanTime")
private String scanTime;
private long scanTime;

public SamzaSqlRelMsgMetadata(@JsonProperty("eventTime") String eventTime, @JsonProperty("arrivalTime") String arrivalTime,
@JsonProperty("scanTime") String scanTime) {
public SamzaSqlRelMsgMetadata(@JsonProperty("eventTime") long eventTime, @JsonProperty("arrivalTime") long arrivalTime,
@JsonProperty("scanTime") long scanTime) {
this.eventTime = eventTime;
this.arrivalTime = arrivalTime;
this.scanTime = scanTime;
}

public SamzaSqlRelMsgMetadata(String eventTime, String arrivalTime, String scanTime, boolean isNewInputMessage) {
public SamzaSqlRelMsgMetadata(long eventTime, long arrivalTime, long scanTime, boolean isNewInputMessage) {
this(eventTime, arrivalTime, scanTime);
this.isNewInputMessage = isNewInputMessage;
}

@JsonProperty("eventTime")
public String getEventTime() { return eventTime;}
public long getEventTime() { return eventTime;}

public void setEventTime(String eventTime) {
public void setEventTime(long eventTime) {
this.eventTime = eventTime;
}

public boolean hasEventTime() { return eventTime != null && !eventTime.isEmpty(); }
public boolean hasEventTime() { return eventTime != 0; }

@JsonProperty("arrivalTime")
public String getarrivalTime() { return arrivalTime;}
public long getarrivalTime() { return arrivalTime;}

public void setArrivalTime(String arrivalTime) {
public void setArrivalTime(long arrivalTime) {
this.arrivalTime = arrivalTime;
}

public boolean hasArrivalTime() { return arrivalTime != null && !arrivalTime.isEmpty(); }
public boolean hasArrivalTime() { return arrivalTime != 0; }

@JsonProperty("scanTime")
public String getscanTime() { return scanTime;}
public long getscanTime() { return scanTime;}

public void setScanTime(String scanTime) {
public void setScanTime(long scanTime) {
this.scanTime = scanTime;
}

public boolean hasScanTime() { return scanTime != null && !scanTime.isEmpty(); }
public boolean hasScanTime() { return scanTime != 0; }

@JsonIgnore
public void setIsSystemMessage(boolean isSystemMessage) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,6 @@

package org.apache.samza.sql.translator;

import com.google.common.annotations.VisibleForTesting;
import java.time.Duration;
import java.time.Instant;
import java.util.Arrays;
import java.util.Collections;
import org.apache.calcite.rel.logical.LogicalFilter;
Expand All @@ -35,6 +32,7 @@
import org.apache.samza.sql.data.Expression;
import org.apache.samza.sql.data.SamzaSqlRelMessage;
import org.apache.samza.sql.runner.SamzaSqlApplicationContext;
import org.apache.samza.sql.util.TranslatorUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -96,7 +94,7 @@ public void init(Context context) {

@Override
public boolean apply(SamzaSqlRelMessage message) {
Instant startProcessing = Instant.now();
long startProcessing = System.nanoTime();
Object[] result = new Object[1];
expr.execute(translatorContext.getExecutionContext(), context, translatorContext.getDataContext(),
message.getSamzaSqlRelRecord().getFieldValues().toArray(), result);
Expand All @@ -105,7 +103,7 @@ public boolean apply(SamzaSqlRelMessage message) {
log.debug(
String.format("return value for input %s is %s",
Arrays.asList(message.getSamzaSqlRelRecord().getFieldValues()).toString(), retVal));
updateMetrics(startProcessing, retVal, Instant.now());
updateMetrics(startProcessing, retVal, System.nanoTime());
return retVal;
} else {
log.error("return value is not boolean");
Expand All @@ -115,17 +113,17 @@ public boolean apply(SamzaSqlRelMessage message) {

/**
* Updates the MetricsRegistery of this operator
* @param startProcessing = begin processing of the message
* @param endProcessing = end of processing
* @param startProcessing = nanoTime when processing of the message started
* @param endProcessing = nanoTIme when processing of the message ended
*/
private void updateMetrics(Instant startProcessing, boolean isOutput, Instant endProcessing) {
private void updateMetrics(long startProcessing, boolean isOutput, long endProcessing) {
inputEvents.inc();
if (isOutput) {
outputEvents.inc();
} else {
filteredOutEvents.inc();
}
processingTime.update(Duration.between(startProcessing, endProcessing).toMillis());
processingTime.update(TranslatorUtils.getMicroDurationFromNanoTimes(startProcessing, endProcessing));
}

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@
package org.apache.samza.sql.translator;

import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import org.apache.calcite.rel.logical.LogicalAggregate;
Expand Down Expand Up @@ -81,7 +80,7 @@ void translate(final LogicalAggregate aggregate, final TranslatorContext context
List<Object> fieldValues = windowPane.getKey().getKey().getSamzaSqlRelRecord().getFieldValues();
fieldNames.add(aggFieldNames.get(0));
fieldValues.add(windowPane.getMessage());
return new SamzaSqlRelMessage(fieldNames, fieldValues, new SamzaSqlRelMsgMetadata("", "", ""));
return new SamzaSqlRelMessage(fieldNames, fieldValues, new SamzaSqlRelMsgMetadata(0, 0,0));
});
context.registerMessageStream(aggregate.getId(), outputStream);
outputStream.map(new TranslatorOutputMetricsMapFunction(logicalOpId));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,6 @@

package org.apache.samza.sql.translator;

import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
Expand All @@ -44,6 +42,7 @@
import org.apache.samza.sql.data.SamzaSqlRelMessage;
import org.apache.samza.sql.data.SamzaSqlRelMsgMetadata;
import org.apache.samza.sql.runner.SamzaSqlApplicationContext;
import org.apache.samza.sql.util.TranslatorUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -113,7 +112,7 @@ public void init(Context context) {
*/
@Override
public SamzaSqlRelMessage apply(SamzaSqlRelMessage message) {
Instant arrivalTime = Instant.now();
long arrivalTime = System.nanoTime();
RelDataType type = project.getRowType();
Object[] output = new Object[type.getFieldCount()];
expr.execute(translatorContext.getExecutionContext(), context, translatorContext.getDataContext(),
Expand All @@ -122,22 +121,22 @@ public SamzaSqlRelMessage apply(SamzaSqlRelMessage message) {
for (int index = 0; index < output.length; index++) {
names.add(index, project.getNamedProjects().get(index).getValue());
}
updateMetrics(arrivalTime, Instant.now(), message.getSamzaSqlRelMsgMetadata().isNewInputMessage);
updateMetrics(arrivalTime, System.nanoTime(), message.getSamzaSqlRelMsgMetadata().isNewInputMessage);
return new SamzaSqlRelMessage(names, Arrays.asList(output), message.getSamzaSqlRelMsgMetadata());
}

/**
* Updates the Diagnostics Metrics (processing time and number of events)
* @param arrivalTime input message arrival time (= beging of processing in this operator)
* @param outputTime output message output time (=end of processing in this operator)
* @param arrivalTime input message arrival nanoTime (= begin of processing in this operator)
* @param outputTime output message output nanoTime (=end of processing in this operator)
* @param isNewInputMessage whether the input Message is from new input message or not
*/
private void updateMetrics(Instant arrivalTime, Instant outputTime, boolean isNewInputMessage) {
private void updateMetrics(long arrivalTime, long outputTime, boolean isNewInputMessage) {
if (isNewInputMessage) {
inputEvents.inc();
}
outputEvents.inc();
processingTime.update(Duration.between(arrivalTime, outputTime).toMillis());
processingTime.update(TranslatorUtils.getMicroDurationFromNanoTimes(arrivalTime, outputTime));
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,6 @@
package org.apache.samza.sql.translator;

import com.google.common.annotations.VisibleForTesting;
import java.time.Duration;
import java.time.Instant;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
Expand Down Expand Up @@ -63,6 +61,7 @@
import org.apache.samza.sql.runner.SamzaSqlApplicationConfig;
import org.apache.samza.sql.runner.SamzaSqlApplicationContext;
import org.apache.samza.sql.util.SamzaSqlQueryParser;
import org.apache.samza.sql.util.TranslatorUtils;
import org.apache.samza.system.descriptors.DelegatingSystemDescriptor;
import org.apache.samza.system.descriptors.GenericOutputDescriptor;
import org.apache.samza.table.Table;
Expand Down Expand Up @@ -134,51 +133,51 @@ public void init(Context context) {

queryLatency = new SamzaHistogram(metricsRegistry, queryLogicalId, TranslatorConstants.QUERY_LATENCY_NAME);
queueingLatency = new SamzaHistogram(metricsRegistry, queryLogicalId, TranslatorConstants.QUEUEING_LATENCY_NAME);

queryOutputEvents = metricsRegistry.newCounter(queryLogicalId, TranslatorConstants.OUTPUT_EVENTS_NAME);
queryOutputEvents.clear();
}

@Override
public KV<Object, Object> apply(SamzaSqlRelMessage message) {
Instant beginProcessing = Instant.now();
long beginProcessing = System.nanoTime();
KV<Object, Object> retKV = this.samzaMsgConverter.convertToSamzaMessage(message);
if (message.getSamzaSqlRelRecord().containsField(SamzaSqlRelMessage.OP_NAME)
&& ((String) message.getSamzaSqlRelRecord().getField(SamzaSqlRelMessage.OP_NAME).get()).equalsIgnoreCase(
SamzaSqlRelMessage.DELETE_OP)) {
// If it is a delete op. Set the payload to null so that the record gets deleted.
retKV = new KV<>(retKV.key, null);
}
updateMetrics(beginProcessing, Instant.now(), message.getSamzaSqlRelMsgMetadata());
updateMetrics(beginProcessing, System.nanoTime(), message.getSamzaSqlRelMsgMetadata());
return retKV;
}

/**
* Updates the Diagnostics Metrics (processing time and number of events)
* @param beginProcessing when sendOutput Started processing this message
* @param endProcessing when sendOutput finished processing this message
* @param beginProcessing when (nanoTime) sendOutput Started processing this message
* @param endProcessing when (nanoTime) sendOutput finished processing this message
* @param metadata the event's message metadata
*/
private void updateMetrics(Instant beginProcessing, Instant endProcessing, SamzaSqlRelMsgMetadata metadata) {
private void updateMetrics(long beginProcessing, long endProcessing, SamzaSqlRelMsgMetadata metadata) {
/* insert (SendToOutputStream) metrics */
insertProcessingTime.update(Duration.between(beginProcessing, endProcessing).toMillis());
insertProcessingTime.update(TranslatorUtils.getMicroDurationFromNanoTimes(beginProcessing, endProcessing));
/* query metrics */
Instant outputTime = Instant.now();
long outputTime = System.nanoTime();
queryOutputEvents.inc();
/* TODO: remove scanTime validation once code to assign it is stable */
Validate.isTrue(metadata.hasScanTime());
Instant scanTime = Instant.parse(metadata.getscanTime());
queryLatency.update(Duration.between(scanTime, outputTime).toMillis());
long scanTime = metadata.getscanTime();
queryLatency.update(TranslatorUtils.getMicroDurationFromNanoTimes(scanTime, outputTime));
/** TODO: change if hasArrivalTime to validation once arrivalTime is assigned,
and later remove the check once code is stable */
if (metadata.hasArrivalTime()) {
Instant arrivalTime = Instant.parse(metadata.getarrivalTime());
queueingLatency.update(Duration.between(arrivalTime, scanTime).toMillis());
long arrivalTime = metadata.getarrivalTime();
queueingLatency.update(TranslatorUtils.getMicroDurationFromNanoTimes(arrivalTime, scanTime));
}
/* since availability of eventTime depends on source, we need the following check */
if (metadata.hasEventTime()) {
Instant eventTime = Instant.parse(metadata.getEventTime());
totalLatency.update(Duration.between(eventTime, outputTime).toMillis());
long eventTime = metadata.getEventTime();
totalLatency.update(TranslatorUtils.getMicroDurationFromNanoTimes(eventTime, outputTime));
}
}
} // OutputMapFunction
Expand Down
Loading