使用Kerberos连接到Mapper内的Accumulo

Gle*_*lic 4 hadoop kerberos accumulo

我正在将一些软件从较旧的Hadoop集群(使用用户名/密码验证)移动到较新的,2.6.0-cdh5.12.0,它启用了Kerberos身份验证.

我已经能够使用AccumuloInput/OutputFormat类中的DelegationToken设置获得许多使用Accumulo的输入和/或输出的现有Map/Reduce作业.

但是,我有1个作业,它使用AccumuloInput/OutputFormat作为输入和输出,但也在它的Mapper.setup()方法中,它通过Zookeeper连接到Accumulo,这样在Mapper.map()方法中,它可以比较每个键/ value正在处理我的Mapper.map()并在另一个Accumulo表中输入.

我在下面列出了相关的代码,其中显示了连接到Zookeeper用户的setup()方法的Password()方法,然后创建了一个Accumulo表Scanner,然后在mapper方法中使用它.

所以问题是我如何使用KerberosToken替换PasswordToken用于在Mapper.setup()方法中设置Accumulo扫描程序?我找不到"获取"我设置的AccumuloInput/OutputFormat类使用的DelegationToken的方法.

我尝试了context.getCredentials().getAllTokens()并查找类型为org.apache.accumulo.code.client.security.tokens.AuthenticationToken的标记 - 此处返回的所有标记都是org.apache.hadoop类型.security.token.Token.

请注意,我输入与剪切/粘贴的代码片段,因为代码在未连接到互联网的网络上运行 - 也就是说可能存在拼写错误.:)

//****************************
// code in the M/R driver
//****************************
ClientConfiguration accumuloCfg = ClientConfiguration.loadDefault().withInstance("Accumulo1").withZkHosts("zookeeper1");
ZooKeeperInstance inst = new ZooKeeperInstance(accumuloCfg);
AuthenticationToken dt = conn.securityOperations().getDelegationToken(new DelagationTokenConfig());
AccumuloInputFormat.setConnectorInfo(job, username, dt);
AccumuloOutputFormat.setConnectorInfo(job, username, dt);
// other job setup and then
job.waitForCompletion(true)



//****************************
// this is inside the Mapper class of the M/R job
//****************************
private Scanner index_scanner;

public void setup(Context context) {
    Configuration cfg = context.getConfiguration();

    // properties set and passed from M/R Driver program
    String username = cfg.get("UserName");
    String password = cfg.get("Password");
    String accumuloInstName = cfg.get("InstanceName");
    String zookeepers = cfg.get("Zookeepers");
    String tableName = cfg.get("TableName");
    Instance inst = new ZooKeeperInstance(accumuloInstName, zookeepers);
    try {
      AuthenticationToken passwordToken = new PasswordToken(password);

      Connector conn = inst.getConnector(username, passwordToken);

      index_scanner = conn.createScanner(tableName, conn.securityOperations().getUserAuthorizations(username));
    } catch(Exception e) {
       e.printStackTrace();
    }
}

public void map(Key key, Value value, Context context) throws IOException, InterruptedException {
    String uuid = key.getRow().toString();
    index_scanner.clearColumns();
    index_scanner.setRange(Range.exact(uuid));
    for(Entry<Key, Value> entry : index_scanner) {
        // do some processing in here
    }
}
Run Code Online (Sandbox Code Playgroud)

Chr*_*her 6

提供的AccumuloInputFormat和AccumuloOutputFormat有一种方法可以在作业配置中设置令牌Accumulo*putFormat.setConnectorInfo(job, principle, token).还可以序列令牌在HDFS文件,使用AuthenticationTokenSerializer和使用的版本,setConnectorInfo它接受一个文件名的方法.

如果传入了KerberosToken,则作业将创建要使用的DelegationToken,如果传入了DelegationToken,则只会使用它.

提供的AccumuloInputFormat应该处理自己的扫描程序,因此通常情况下,如果您正确设置了配置,则不必在Mapper中执行此操作.但是,如果您在Mapper中进行二次扫描(对于类似连接的东西),您可以检查提供AccumuloInputFormat的RecordReader源代码,以获取如何检索配置和构建扫描器的示例.