mirror of https://github.com/apache/druid.git
1) Serialization of Tasks is important
This commit is contained in:
parent
f1175389c4
commit
688e5e7417
|
@ -111,4 +111,32 @@ public abstract class AbstractTask implements Task
|
||||||
{
|
{
|
||||||
return TaskStatus.success(getId());
|
return TaskStatus.success(getId());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean equals(Object o)
|
||||||
|
{
|
||||||
|
if (this == o) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
if (o == null || getClass() != o.getClass()) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
AbstractTask that = (AbstractTask) o;
|
||||||
|
|
||||||
|
if (dataSource != null ? !dataSource.equals(that.dataSource) : that.dataSource != null) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
if (groupId != null ? !groupId.equals(that.groupId) : that.groupId != null) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
if (id != null ? !id.equals(that.id) : that.id != null) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
if (interval != null ? !interval.equals(that.interval) : that.interval != null) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
return true;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -19,6 +19,7 @@
|
||||||
|
|
||||||
package com.metamx.druid.merger.common.task;
|
package com.metamx.druid.merger.common.task;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.annotation.JsonCreator;
|
||||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||||
import com.google.common.base.Function;
|
import com.google.common.base.Function;
|
||||||
import com.google.common.collect.Lists;
|
import com.google.common.collect.Lists;
|
||||||
|
@ -53,14 +54,37 @@ public class VersionConverterTask extends AbstractTask
|
||||||
private static final Logger log = new Logger(VersionConverterTask.class);
|
private static final Logger log = new Logger(VersionConverterTask.class);
|
||||||
private final DataSegment segment;
|
private final DataSegment segment;
|
||||||
|
|
||||||
public VersionConverterTask(
|
public static VersionConverterTask create(String dataSource, Interval interval)
|
||||||
|
{
|
||||||
|
final String id = makeId(dataSource, interval);
|
||||||
|
return new VersionConverterTask(id, id, dataSource, interval, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
public static VersionConverterTask create(DataSegment segment)
|
||||||
|
{
|
||||||
|
final Interval interval = segment.getInterval();
|
||||||
|
final String dataSource = segment.getDataSource();
|
||||||
|
final String id = makeId(dataSource, interval);
|
||||||
|
return new VersionConverterTask(id, id, dataSource, interval, segment);
|
||||||
|
}
|
||||||
|
|
||||||
|
private static String makeId(String dataSource, Interval interval)
|
||||||
|
{
|
||||||
|
return joinId(TYPE, dataSource, interval.getStart(), interval.getEnd(), new DateTime());
|
||||||
|
}
|
||||||
|
|
||||||
|
@JsonCreator
|
||||||
|
private VersionConverterTask(
|
||||||
|
@JsonProperty("id") String id,
|
||||||
|
@JsonProperty("groupId") String groupId,
|
||||||
@JsonProperty("dataSource") String dataSource,
|
@JsonProperty("dataSource") String dataSource,
|
||||||
@JsonProperty("interval") Interval interval,
|
@JsonProperty("interval") Interval interval,
|
||||||
@JsonProperty("segment") DataSegment segment
|
@JsonProperty("segment") DataSegment segment
|
||||||
)
|
)
|
||||||
{
|
{
|
||||||
super(
|
super(
|
||||||
joinId(TYPE, dataSource, interval.getStart(), interval.getEnd(), new DateTime()),
|
id,
|
||||||
|
groupId,
|
||||||
dataSource,
|
dataSource,
|
||||||
interval
|
interval
|
||||||
);
|
);
|
||||||
|
@ -74,6 +98,12 @@ public class VersionConverterTask extends AbstractTask
|
||||||
return TYPE;
|
return TYPE;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@JsonProperty
|
||||||
|
public DataSegment getSegment()
|
||||||
|
{
|
||||||
|
return segment;
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public TaskStatus run(TaskToolbox toolbox) throws Exception
|
public TaskStatus run(TaskToolbox toolbox) throws Exception
|
||||||
{
|
{
|
||||||
|
@ -121,11 +151,31 @@ public class VersionConverterTask extends AbstractTask
|
||||||
return TaskStatus.success(getId());
|
return TaskStatus.success(getId());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public boolean equals(Object o)
|
||||||
|
{
|
||||||
|
if (this == o) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
if (o == null || getClass() != o.getClass()) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
VersionConverterTask that = (VersionConverterTask) o;
|
||||||
|
|
||||||
|
if (segment != null ? !segment.equals(that.segment) : that.segment != null) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
return super.equals(o);
|
||||||
|
}
|
||||||
|
|
||||||
public static class SubTask extends AbstractTask
|
public static class SubTask extends AbstractTask
|
||||||
{
|
{
|
||||||
private final DataSegment segment;
|
private final DataSegment segment;
|
||||||
|
|
||||||
protected SubTask(
|
@JsonCreator
|
||||||
|
public SubTask(
|
||||||
@JsonProperty("groupId") String groupId,
|
@JsonProperty("groupId") String groupId,
|
||||||
@JsonProperty("segment") DataSegment segment
|
@JsonProperty("segment") DataSegment segment
|
||||||
)
|
)
|
||||||
|
@ -145,6 +195,12 @@ public class VersionConverterTask extends AbstractTask
|
||||||
this.segment = segment;
|
this.segment = segment;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@JsonProperty
|
||||||
|
public DataSegment getSegment()
|
||||||
|
{
|
||||||
|
return segment;
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String getType()
|
public String getType()
|
||||||
{
|
{
|
||||||
|
|
|
@ -0,0 +1,66 @@
|
||||||
|
/*
|
||||||
|
* Druid - a distributed column store.
|
||||||
|
* Copyright (C) 2012 Metamarkets Group Inc.
|
||||||
|
*
|
||||||
|
* This program is free software; you can redistribute it and/or
|
||||||
|
* modify it under the terms of the GNU General Public License
|
||||||
|
* as published by the Free Software Foundation; either version 2
|
||||||
|
* of the License, or (at your option) any later version.
|
||||||
|
*
|
||||||
|
* This program is distributed in the hope that it will be useful,
|
||||||
|
* but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
* GNU General Public License for more details.
|
||||||
|
*
|
||||||
|
* You should have received a copy of the GNU General Public License
|
||||||
|
* along with this program; if not, write to the Free Software
|
||||||
|
* Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301, USA.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package com.metamx.druid.merger.common.task;
|
||||||
|
|
||||||
|
import com.google.common.collect.ImmutableList;
|
||||||
|
import com.google.common.collect.ImmutableMap;
|
||||||
|
import com.metamx.druid.client.DataSegment;
|
||||||
|
import com.metamx.druid.jackson.DefaultObjectMapper;
|
||||||
|
import com.metamx.druid.shard.NoneShardSpec;
|
||||||
|
import junit.framework.Assert;
|
||||||
|
import org.joda.time.DateTime;
|
||||||
|
import org.joda.time.Interval;
|
||||||
|
import org.junit.Test;
|
||||||
|
|
||||||
|
/**
|
||||||
|
*/
|
||||||
|
public class VersionConverterTaskTest
|
||||||
|
{
|
||||||
|
@Test
|
||||||
|
public void testSerializationSimple() throws Exception
|
||||||
|
{
|
||||||
|
final String dataSource = "billy";
|
||||||
|
final Interval interval = new Interval(new DateTime().minus(1000), new DateTime());
|
||||||
|
|
||||||
|
DefaultObjectMapper jsonMapper = new DefaultObjectMapper();
|
||||||
|
|
||||||
|
VersionConverterTask task = VersionConverterTask.create(dataSource, interval);
|
||||||
|
|
||||||
|
Task task2 = jsonMapper.readValue(jsonMapper.writerWithDefaultPrettyPrinter().writeValueAsString(task), Task.class);
|
||||||
|
Assert.assertEquals(task, task2);
|
||||||
|
|
||||||
|
DataSegment segment = new DataSegment(
|
||||||
|
dataSource,
|
||||||
|
interval,
|
||||||
|
new DateTime().toString(),
|
||||||
|
ImmutableMap.<String, Object>of(),
|
||||||
|
ImmutableList.<String>of(),
|
||||||
|
ImmutableList.<String>of(),
|
||||||
|
new NoneShardSpec(),
|
||||||
|
9,
|
||||||
|
102937
|
||||||
|
);
|
||||||
|
|
||||||
|
task = VersionConverterTask.create(segment);
|
||||||
|
|
||||||
|
task2 = jsonMapper.readValue(jsonMapper.writerWithDefaultPrettyPrinter().writeValueAsString(task), Task.class);
|
||||||
|
Assert.assertEquals(task, task2);
|
||||||
|
}
|
||||||
|
}
|
Loading…
Reference in New Issue