Skip to content

Commit cf37442

Browse files
committed
add raw event API
Signed-off-by: Maciej Obuchowski <obuchowski.maciej@gmail.com>
1 parent 7c6ddd8 commit cf37442

19 files changed

Lines changed: 450 additions & 1 deletion

api/build.gradle

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ dependencies {
4040
implementation "io.prometheus:simpleclient_hotspot:${prometheusVersion}"
4141
implementation "io.prometheus:simpleclient_servlet:${prometheusVersion}"
4242
implementation "org.jdbi:jdbi3-core:${jdbi3Version}"
43+
implementation "org.jdbi:jdbi3-jackson2:${jdbi3Version}"
4344
implementation "org.jdbi:jdbi3-postgres:${jdbi3Version}"
4445
implementation "org.jdbi:jdbi3-sqlobject:${jdbi3Version}"
4546
implementation 'com.google.guava:guava:31.1-jre'

api/src/main/java/marquez/MarquezContext.java

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
import lombok.Getter;
1515
import lombok.NonNull;
1616
import marquez.api.DatasetResource;
17+
import marquez.api.EventResource;
1718
import marquez.api.JobResource;
1819
import marquez.api.NamespaceResource;
1920
import marquez.api.OpenLineageResource;
@@ -25,6 +26,7 @@
2526
import marquez.db.DatasetDao;
2627
import marquez.db.DatasetFieldDao;
2728
import marquez.db.DatasetVersionDao;
29+
import marquez.db.EventDao;
2830
import marquez.db.JobContextDao;
2931
import marquez.db.JobDao;
3032
import marquez.db.JobVersionDao;
@@ -42,6 +44,7 @@
4244
import marquez.service.DatasetFieldService;
4345
import marquez.service.DatasetService;
4446
import marquez.service.DatasetVersionService;
47+
import marquez.service.EventService;
4548
import marquez.service.JobService;
4649
import marquez.service.LineageService;
4750
import marquez.service.NamespaceService;
@@ -71,6 +74,7 @@ public final class MarquezContext {
7174
@Getter private final OpenLineageDao openLineageDao;
7275
@Getter private final LineageDao lineageDao;
7376
@Getter private final SearchDao searchDao;
77+
@Getter private final EventDao eventDao;
7478

7579
@Getter private final List<RunTransitionListener> runTransitionListeners;
7680

@@ -82,6 +86,7 @@ public final class MarquezContext {
8286
@Getter private final RunService runService;
8387
@Getter private final OpenLineageService openLineageService;
8488
@Getter private final LineageService lineageService;
89+
@Getter private final EventService eventService;
8590

8691
@Getter private final NamespaceResource namespaceResource;
8792
@Getter private final SourceResource sourceResource;
@@ -90,6 +95,7 @@ public final class MarquezContext {
9095
@Getter private final TagResource tagResource;
9196
@Getter private final OpenLineageResource openLineageResource;
9297
@Getter private final SearchResource searchResource;
98+
@Getter private final EventResource eventResource;
9399

94100
@Getter private final ImmutableList<Object> resources;
95101
@Getter private final JdbiExceptionExceptionMapper jdbiException;
@@ -119,6 +125,7 @@ private MarquezContext(
119125
this.openLineageDao = jdbi.onDemand(OpenLineageDao.class);
120126
this.lineageDao = jdbi.onDemand(LineageDao.class);
121127
this.searchDao = jdbi.onDemand(SearchDao.class);
128+
this.eventDao = jdbi.onDemand(EventDao.class);
122129
this.runTransitionListeners = runTransitionListeners;
123130

124131
this.namespaceService = new NamespaceService(baseDao);
@@ -131,6 +138,7 @@ private MarquezContext(
131138
this.tagService.init(tags);
132139
this.openLineageService = new OpenLineageService(baseDao, runService);
133140
this.lineageService = new LineageService(lineageDao, jobDao);
141+
this.eventService = new EventService(baseDao);
134142
this.jdbiException = new JdbiExceptionExceptionMapper();
135143
final ServiceFactory serviceFactory =
136144
ServiceFactory.builder()
@@ -144,6 +152,7 @@ private MarquezContext(
144152
.lineageService(lineageService)
145153
.datasetFieldService(new DatasetFieldService(baseDao))
146154
.datasetVersionService(new DatasetVersionService(baseDao))
155+
.eventService(eventService)
147156
.build();
148157
this.namespaceResource = new NamespaceResource(serviceFactory);
149158
this.sourceResource = new SourceResource(serviceFactory);
@@ -152,6 +161,7 @@ private MarquezContext(
152161
this.tagResource = new TagResource(serviceFactory);
153162
this.openLineageResource = new OpenLineageResource(serviceFactory);
154163
this.searchResource = new SearchResource(searchDao);
164+
this.eventResource = new EventResource(serviceFactory);
155165

156166
this.resources =
157167
ImmutableList.of(
@@ -162,7 +172,8 @@ private MarquezContext(
162172
tagResource,
163173
jdbiException,
164174
openLineageResource,
165-
searchResource);
175+
searchResource,
176+
eventResource);
166177

167178
final MarquezGraphqlServletBuilder servlet = new MarquezGraphqlServletBuilder();
168179
this.graphqlServlet = servlet.getServlet(new GraphqlSchemaBuilder(jdbi));

api/src/main/java/marquez/api/BaseResource.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import marquez.service.DatasetFieldService;
2929
import marquez.service.DatasetService;
3030
import marquez.service.DatasetVersionService;
31+
import marquez.service.EventService;
3132
import marquez.service.JobService;
3233
import marquez.service.LineageService;
3334
import marquez.service.NamespaceService;
@@ -50,6 +51,7 @@ public class BaseResource {
5051
protected DatasetVersionService datasetVersionService;
5152
protected DatasetFieldService datasetFieldService;
5253
protected LineageService lineageService;
54+
protected EventService eventService;
5355

5456
public BaseResource(ServiceFactory serviceFactory) {
5557
this.serviceFactory = serviceFactory;
@@ -63,6 +65,7 @@ public BaseResource(ServiceFactory serviceFactory) {
6365
this.datasetVersionService = serviceFactory.getDatasetVersionService();
6466
this.datasetFieldService = serviceFactory.getDatasetFieldService();
6567
this.lineageService = serviceFactory.getLineageService();
68+
this.eventService = serviceFactory.getEventService();
6669
}
6770

6871
void throwIfNotExists(@NonNull NamespaceName namespaceName) {
Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,75 @@
1+
package marquez.api;
2+
3+
import com.codahale.metrics.annotation.ExceptionMetered;
4+
import com.codahale.metrics.annotation.ResponseMetered;
5+
import com.codahale.metrics.annotation.Timed;
6+
import com.fasterxml.jackson.annotation.JsonProperty;
7+
import lombok.NonNull;
8+
import lombok.Value;
9+
import marquez.api.exceptions.NamespaceNotFoundException;
10+
import marquez.common.models.NamespaceName;
11+
import marquez.service.ServiceFactory;
12+
import marquez.service.models.LineageEvent;
13+
import marquez.service.models.Namespace;
14+
import marquez.service.models.NamespaceMeta;
15+
import marquez.service.models.RawLineageEvent;
16+
17+
import javax.validation.Valid;
18+
import javax.validation.constraints.Min;
19+
import javax.ws.rs.Consumes;
20+
import javax.ws.rs.DefaultValue;
21+
import javax.ws.rs.GET;
22+
import javax.ws.rs.PUT;
23+
import javax.ws.rs.Path;
24+
import javax.ws.rs.PathParam;
25+
import javax.ws.rs.Produces;
26+
import javax.ws.rs.QueryParam;
27+
import javax.ws.rs.core.Response;
28+
29+
import java.util.List;
30+
31+
import static javax.ws.rs.core.MediaType.APPLICATION_JSON;
32+
33+
@Path("/api/v1")
34+
public class EventResource extends BaseResource {
35+
public EventResource(@NonNull final ServiceFactory serviceFactory) {
36+
super(serviceFactory);
37+
}
38+
39+
@Timed
40+
@ResponseMetered
41+
@ExceptionMetered
42+
@GET
43+
@Path("/events")
44+
@Produces(APPLICATION_JSON)
45+
public Response get(
46+
@QueryParam("limit") @DefaultValue("100") @Min(value = 0) int limit,
47+
@QueryParam("offset") @DefaultValue("0") @Min(value = 0) int offset) {
48+
final List<RawLineageEvent> events =
49+
eventService.getAll(limit, offset);
50+
return Response.ok(new Events(events)).build();
51+
}
52+
53+
@Timed
54+
@ResponseMetered
55+
@ExceptionMetered
56+
@GET
57+
@Path("/events/{namespace}")
58+
@Produces(APPLICATION_JSON)
59+
public Response getByNamespace(
60+
@PathParam("namespace") String namespace,
61+
@QueryParam("limit") @DefaultValue("100") @Min(value = 0) int limit,
62+
@QueryParam("offset") @DefaultValue("0") @Min(value = 0) int offset) {
63+
final List<RawLineageEvent> event =
64+
eventService.getByNamespace(namespace, limit, offset);
65+
return Response.ok(new Events(event)).build();
66+
}
67+
68+
@Value
69+
static class Events {
70+
@NonNull
71+
@JsonProperty("events")
72+
List<RawLineageEvent> value;
73+
}
74+
75+
}

api/src/main/java/marquez/db/BaseDao.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,4 +50,7 @@ public interface BaseDao extends SqlObject {
5050

5151
@CreateSqlObject
5252
OpenLineageDao createOpenLineageDao();
53+
54+
@CreateSqlObject
55+
EventDao createEventDao();
5356
}

api/src/main/java/marquez/db/Columns.java

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@
1717
import java.sql.ResultSetMetaData;
1818
import java.sql.SQLException;
1919
import java.time.Instant;
20+
import java.time.ZoneId;
21+
import java.time.ZonedDateTime;
2022
import java.util.Arrays;
2123
import java.util.Collections;
2224
import java.util.List;
@@ -126,6 +128,13 @@ private Columns() {}
126128
public static final String RUN_UUID = "run_uuid";
127129
public static final String STATE = "state";
128130

131+
/* LINEAGE EVENT ROW COLUMNS */
132+
133+
public static final String EVENT_TIME = "event_time";
134+
public static final String EVENT = "event";
135+
public static final String EVENT_TYPE = "event_type";
136+
public static final String PRODUCER = "producer";
137+
129138
public static UUID uuidOrNull(final ResultSet results, final String column) throws SQLException {
130139
if (results.getObject(column) == null) {
131140
return null;
@@ -156,6 +165,15 @@ public static Instant timestampOrThrow(final ResultSet results, final String col
156165
return results.getTimestamp(column).toInstant();
157166
}
158167

168+
public static ZonedDateTime zonedDateTimeOrThrow(final ResultSet results, final String column)
169+
throws SQLException {
170+
if (results.getObject(column) == null) {
171+
throw new IllegalArgumentException();
172+
}
173+
return results.getTimestamp(column).toInstant().atZone(ZoneId.of("UTC"));
174+
}
175+
176+
159177
public static String stringOrNull(final ResultSet results, final String column)
160178
throws SQLException {
161179
if (results.getObject(column) == null) {
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
package marquez.db;
2+
3+
4+
import marquez.db.mappers.RawLineageEventMapper;
5+
import marquez.service.models.RawLineageEvent;
6+
import org.jdbi.v3.sqlobject.config.RegisterRowMapper;
7+
import org.jdbi.v3.sqlobject.customizer.Bind;
8+
import org.jdbi.v3.sqlobject.statement.SqlQuery;
9+
import java.util.List;
10+
11+
@RegisterRowMapper(RawLineageEventMapper.class)
12+
public interface EventDao extends BaseDao {
13+
14+
@SqlQuery("""
15+
SELECT *
16+
FROM lineage_events
17+
ORDER BY event_time DESC
18+
LIMIT :limit
19+
OFFSET :offset""")
20+
List<RawLineageEvent> getAll(int limit, int offset);
21+
22+
23+
/**
24+
* This is a "hack" to get inputs/outputs namespace from jsonb column:
25+
* <a href="https://github.com/jdbi/jdbi/issues/1510#issuecomment-485423083">explanation</a>
26+
*/
27+
@SqlQuery("""
28+
SELECT le.*
29+
FROM lineage_events le, jsonb_array_elements(coalesce(le.event -> 'inputs', '[]'::jsonb) || coalesce(le.event -> 'outputs', '[]'::jsonb)) AS ds
30+
WHERE le.job_namespace = :namespace
31+
OR ds ->> 'namespace' = :namespace
32+
LIMIT :limit
33+
OFFSET :offset""")
34+
List<RawLineageEvent> getByNamespace(@Bind("namespace") String namespace, @Bind("limit") int limit, @Bind("offset") int offset);
35+
}
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
package marquez.db.mappers;
2+
3+
import com.fasterxml.jackson.core.JsonProcessingException;
4+
import com.fasterxml.jackson.databind.ObjectMapper;
5+
import lombok.extern.slf4j.Slf4j;
6+
import marquez.common.Utils;
7+
import marquez.db.Columns;
8+
import marquez.service.models.LineageEvent;
9+
import marquez.service.models.RawLineageEvent;
10+
import org.jdbi.v3.core.mapper.RowMapper;
11+
import org.jdbi.v3.core.statement.StatementContext;
12+
13+
import java.sql.ResultSet;
14+
import java.sql.SQLException;
15+
import java.util.Collections;
16+
import java.util.List;
17+
import java.util.Map;
18+
19+
import static marquez.db.Columns.stringOrThrow;
20+
import static marquez.db.Columns.zonedDateTimeOrThrow;
21+
22+
@Slf4j
23+
public class RawLineageEventMapper implements RowMapper<RawLineageEvent> {
24+
@Override
25+
public RawLineageEvent map(ResultSet rs, StatementContext ctx) throws SQLException {
26+
String rawEvent = stringOrThrow(rs, Columns.EVENT);
27+
Map<String, Object> run = Collections.emptyMap();
28+
Map<String, Object> job = Collections.emptyMap();
29+
30+
List<Object> inputs = Collections.emptyList();
31+
List<Object> outputs = Collections.emptyList();
32+
try {
33+
ObjectMapper mapper = Utils.getMapper();
34+
Map<String, Object> event = mapper.readValue(rawEvent, Map.class);
35+
run = (Map<String, Object>) event.getOrDefault("run", Collections.emptyMap());
36+
job = (Map<String, Object>) event.getOrDefault("job", Collections.emptyMap());
37+
inputs = (List<Object>) event.getOrDefault("inputs", Collections.emptyList());
38+
outputs = (List<Object>) event.getOrDefault("outputs", Collections.emptyList());
39+
} catch (JsonProcessingException e) {
40+
log.error("Failed to process json", e);
41+
}
42+
43+
return RawLineageEvent.builder()
44+
.eventTime(zonedDateTimeOrThrow(rs, Columns.EVENT_TIME))
45+
.eventType(stringOrThrow(rs, Columns.EVENT_TYPE))
46+
.run(run)
47+
.job(job)
48+
.inputs(inputs)
49+
.outputs(outputs)
50+
.producer(stringOrThrow(rs, Columns.PRODUCER))
51+
.build();
52+
}
53+
}

api/src/main/java/marquez/service/DelegatingDaos.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
import marquez.db.DatasetDao;
1111
import marquez.db.DatasetFieldDao;
1212
import marquez.db.DatasetVersionDao;
13+
import marquez.db.EventDao;
1314
import marquez.db.JobContextDao;
1415
import marquez.db.JobDao;
1516
import marquez.db.JobVersionDao;
@@ -98,4 +99,10 @@ public static class DelegatingTagDao implements TagDao {
9899
public static class DelegatingLineageDao implements LineageDao {
99100
@Delegate private final LineageDao delegate;
100101
}
102+
103+
@RequiredArgsConstructor
104+
public static class DelegatingEventDao implements EventDao {
105+
@Delegate private final EventDao delegate;
106+
}
107+
101108
}
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
package marquez.service;
2+
3+
import lombok.NonNull;
4+
import marquez.db.BaseDao;
5+
6+
public class EventService extends DelegatingDaos.DelegatingEventDao {
7+
public EventService(@NonNull BaseDao baseDao) {
8+
super(baseDao.createEventDao());
9+
}
10+
}

0 commit comments

Comments
 (0)