Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Issue 733: Event Time support in java client #734

Merged
merged 1 commit into from
Aug 31, 2017
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 @@ -73,8 +73,25 @@ public interface Message {
*/
MessageId getMessageId();

/**
* Get the publish time of this message. The publish time is the timestamp that a client publish the message.
*
* @return publish time of this message.
* @see #getEventTime()
*/
long getPublishTime();

/**
* Get the event time associated with this message. It is typically set by the applications via
* {@link MessageBuilder#setEventTime(long)}.
*
* <p>If there isn't any event time associated with this event, it will return 0.
*
* @see MessageBuilder#setEventTime(long)
* @since 1.20.0
*/
long getEventTime();

/**
* Check whether the message has a key
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,18 @@ public static MessageBuilder create() {
*/
MessageBuilder setKey(String key);

/**
* Set the event time for a given message.
*
* <p>Applications can retrieve the event time by calling {@link Message#getEventTime()}.
*
* <p>Note: currently pulsar doesn't support event-time based index. so the subscribers can't
* seek the messages by event time.
*
* @since 1.20.0
*/
MessageBuilder setEventTime(long timestamp);

/**
* Override the replication clusters for this message.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
*/
package org.apache.pulsar.client.impl;

import static com.google.common.base.Preconditions.checkArgument;

import java.nio.ByteBuffer;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -79,6 +81,13 @@ public MessageBuilder setKey(String key) {
return this;
}

@Override
public MessageBuilder setEventTime(long timestamp) {
checkArgument(timestamp > 0, "Invalid timestamp : '%s'", timestamp);
msgMetadataBuilder.setEventTime(timestamp);
return this;
}

@Override
public MessageBuilder setReplicationClusters(List<String> clusters) {
Preconditions.checkNotNull(clusters);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,15 @@ public long getPublishTime() {
return msgMetadataBuilder.getPublishTime();
}

@Override
public long getEventTime() {
checkNotNull(msgMetadataBuilder);
if (msgMetadataBuilder.hasEventTime()) {
return msgMetadataBuilder.getEventTime();
}
return 0;
}

public boolean isExpired(int messageTTLInSeconds) {
return messageTTLInSeconds != 0
&& System.currentTimeMillis() > (getPublishTime() + TimeUnit.SECONDS.toMillis(messageTTLInSeconds));
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
/**
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.pulsar.client.impl;

import static org.junit.Assert.assertEquals;

import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageBuilder;
import org.junit.Test;

/**
* Unit test of {@link MessageBuilderImpl}.
*/
public class MessageBuilderTest {

@Test(expected = IllegalArgumentException.class)
public void testSetEventTimeNegative() {
MessageBuilder builder = MessageBuilder.create();
builder.setEventTime(-1L);
}

@Test(expected = IllegalArgumentException.class)
public void testSetEventTimeZero() {
MessageBuilder builder = MessageBuilder.create();
builder.setEventTime(0L);
}

@Test
public void testSetEventTimePositive() {
long eventTime = System.currentTimeMillis();
MessageBuilder builder = MessageBuilder.create();
builder.setContent(new byte[0]);
builder.setEventTime(eventTime);
Message msg = builder.build();
assertEquals(eventTime, msg.getEventTime());
}

@Test
public void testBuildMessageWithoutEventTime() {
MessageBuilder builder = MessageBuilder.create();
builder.setContent(new byte[0]);
Message msg = builder.build();
assertEquals(0L, msg.getEventTime());
}

}

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 4 additions & 0 deletions pulsar-common/src/main/proto/PulsarApi.proto
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,10 @@ message MessageMetadata {
//optional sfixed64 checksum = 10;
// differentiate single and batch message metadata
optional int32 num_messages_in_batch = 11 [default = 1];

// the timestamp that this event occurs. it is typically set by applications.
// if this field is omitted, `publish_time` can be used for the purpose of `event_time`.
optional uint64 event_time = 12 [default = 0];
}


Expand Down