public class org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema extends org.apache.drill.exec.store.AbstractSchema
{
private static final org.slf4j.Logger logger;
private final org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin;
private final java.util.Map drillTables;
private java.util.Set tableNames;
public void <init>(org.apache.drill.exec.store.kafka.KafkaStoragePlugin, java.lang.String)
{
org.apache.drill.exec.store.kafka.KafkaStoragePlugin v;
java.util.List v;
org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema v;
java.util.HashMap v;
java.lang.String v;
v := @this: org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema;
v := @parameter: org.apache.drill.exec.store.kafka.KafkaStoragePlugin;
v := @parameter: java.lang.String;
v = staticinvoke <java.util.Collections: java.util.List emptyList()>();
specialinvoke v.<org.apache.drill.exec.store.AbstractSchema: void <init>(java.util.List,java.lang.String)>(v, v);
v = new java.util.HashMap;
specialinvoke v.<java.util.HashMap: void <init>()>();
v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Map drillTables> = v;
v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin> = v;
return;
}
public java.lang.String getTypeName()
{
org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema v;
v := @this: org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema;
return "kafka";
}
void setHolder(org.apache.calcite.schema.SchemaPlus)
{
org.apache.calcite.schema.SchemaPlus v;
java.util.Iterator v;
java.util.Set v;
org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema v;
java.lang.Object v;
org.apache.calcite.schema.Schema v;
boolean v;
v := @this: org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema;
v := @parameter: org.apache.calcite.schema.SchemaPlus;
v = virtualinvoke v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Set getSubSchemaNames()>();
v = interfaceinvoke v.<java.util.Set: java.util.Iterator iterator()>();
label:
v = interfaceinvoke v.<java.util.Iterator: boolean hasNext()>();
if v == 0 goto label;
v = interfaceinvoke v.<java.util.Iterator: java.lang.Object next()>();
v = virtualinvoke v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.calcite.schema.Schema getSubSchema(java.lang.String)>(v);
interfaceinvoke v.<org.apache.calcite.schema.SchemaPlus: org.apache.calcite.schema.SchemaPlus add(java.lang.String,org.apache.calcite.schema.Schema)>(v, v);
goto label;
label:
return;
}
public org.apache.calcite.schema.Table getTable(java.lang.String)
{
org.apache.drill.exec.store.kafka.KafkaScanSpec v;
org.apache.drill.exec.planner.logical.DynamicDrillTable v;
org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema v;
org.apache.drill.exec.store.kafka.KafkaStoragePlugin v;
java.util.Map v, v, v;
java.lang.Object v;
java.lang.String v, v;
boolean v;
v := @this: org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema;
v := @parameter: java.lang.String;
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Map drillTables>;
v = interfaceinvoke v.<java.util.Map: boolean containsKey(java.lang.Object)>(v);
if v != 0 goto label;
v = new org.apache.drill.exec.store.kafka.KafkaScanSpec;
specialinvoke v.<org.apache.drill.exec.store.kafka.KafkaScanSpec: void <init>(java.lang.String)>(v);
v = new org.apache.drill.exec.planner.logical.DynamicDrillTable;
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin>;
v = virtualinvoke v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.lang.String getName()>();
specialinvoke v.<org.apache.drill.exec.planner.logical.DynamicDrillTable: void <init>(org.apache.drill.exec.store.StoragePlugin,java.lang.String,org.apache.drill.exec.planner.logical.DrillTableSelection)>(v, v, v);
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Map drillTables>;
interfaceinvoke v.<java.util.Map: java.lang.Object put(java.lang.Object,java.lang.Object)>(v, v);
label:
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Map drillTables>;
v = interfaceinvoke v.<java.util.Map: java.lang.Object get(java.lang.Object)>(v);
return v;
}
public java.util.Set getTableNames()
{
java.lang.Throwable v, v;
java.lang.Object[] v;
java.util.Map v;
java.lang.String v, v;
java.util.Properties v;
org.slf4j.Logger v;
java.util.Set v, v, v, v;
org.apache.drill.exec.store.kafka.KafkaStoragePluginConfig v;
org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema v;
java.lang.Exception v;
org.apache.drill.exec.store.kafka.KafkaStoragePlugin v, v, v, v;
org.apache.kafka.clients.consumer.KafkaConsumer v, v;
v := @this: org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema;
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Set tableNames>;
if v != null goto label;
v = null;
label:
v = new org.apache.kafka.clients.consumer.KafkaConsumer;
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin>;
v = virtualinvoke v.<org.apache.drill.exec.store.kafka.KafkaStoragePlugin: org.apache.drill.exec.store.kafka.KafkaStoragePluginConfig getConfig()>();
v = virtualinvoke v.<org.apache.drill.exec.store.kafka.KafkaStoragePluginConfig: java.util.Properties getKafkaConsumerProps()>();
specialinvoke v.<org.apache.kafka.clients.consumer.KafkaConsumer: void <init>(java.util.Properties)>(v);
v = v;
v = virtualinvoke v.<org.apache.kafka.clients.consumer.KafkaConsumer: java.util.Map listTopics()>();
v = interfaceinvoke v.<java.util.Map: java.util.Set keySet()>();
v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Set tableNames> = v;
label:
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin>;
virtualinvoke v.<org.apache.drill.exec.store.kafka.KafkaStoragePlugin: void registerToClose(java.lang.AutoCloseable)>(v);
goto label;
label:
v := @caughtexception;
v = <org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.slf4j.Logger logger>;
v = newarray (java.lang.Object)[3];
v = virtualinvoke v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.lang.String getName()>();
v[0] = v;
v = virtualinvoke v.<java.lang.Exception: java.lang.String getMessage()>();
v[1] = v;
v = virtualinvoke v.<java.lang.Exception: java.lang.Throwable getCause()>();
v[2] = v;
interfaceinvoke v.<org.slf4j.Logger: void warn(java.lang.String,java.lang.Object[])>("Failure while loading table names for database \'{}\': {}", v);
v = staticinvoke <java.util.Collections: java.util.Set emptySet()>();
label:
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin>;
virtualinvoke v.<org.apache.drill.exec.store.kafka.KafkaStoragePlugin: void registerToClose(java.lang.AutoCloseable)>(v);
return v;
label:
v := @caughtexception;
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.apache.drill.exec.store.kafka.KafkaStoragePlugin plugin>;
virtualinvoke v.<org.apache.drill.exec.store.kafka.KafkaStoragePlugin: void registerToClose(java.lang.AutoCloseable)>(v);
throw v;
label:
v = v.<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: java.util.Set tableNames>;
return v;
catch java.lang.Exception from label to label with label;
catch java.lang.Throwable from label to label with label;
catch java.lang.Throwable from label to label with label;
}
static void <clinit>()
{
org.slf4j.Logger v;
v = staticinvoke <org.slf4j.LoggerFactory: org.slf4j.Logger getLogger(java.lang.Class)>(class "Lorg/apache/drill/exec/store/kafka/schema/KafkaMessageSchema;");
<org.apache.drill.exec.store.kafka.schema.KafkaMessageSchema: org.slf4j.Logger logger> = v;
return;
}
}